Class KafkaProducerMessageHandler<K,V>
java.lang.Object
org.springframework.integration.context.IntegrationObjectSupport
org.springframework.integration.handler.MessageHandlerSupport
org.springframework.integration.handler.AbstractMessageHandler
org.springframework.integration.handler.AbstractMessageProducingHandler
org.springframework.integration.handler.AbstractReplyProducingMessageHandler
org.springframework.integration.kafka.outbound.KafkaProducerMessageHandler<K,V>
- Type Parameters:
K- the key type.V- the value type.
- All Implemented Interfaces:
org.reactivestreams.Subscriber<Message<?>>,Aware,BeanClassLoaderAware,BeanFactoryAware,BeanNameAware,DisposableBean,InitializingBean,ApplicationContextAware,Lifecycle,Ordered,ExpressionCapable,Orderable,MessageProducer,HeaderPropagationAware,IntegrationPattern,NamedComponent,IntegrationManagement,ManageableLifecycle,TrackableComponent,MessageHandler,reactor.core.CoreSubscriber<Message<?>>
public class KafkaProducerMessageHandler<K,V> extends AbstractReplyProducingMessageHandler implements ManageableLifecycle
Kafka Message Handler; when supplied with a
ReplyingKafkaTemplate it is used as
the handler in an outbound gateway. When supplied with a simple KafkaTemplate
it used as the handler in an outbound channel adapter.
Starting with version 3.2.1 the handler supports receiving a pre-built
ProducerRecord payload. In that case, most configuration properties
(setTopicExpression(Expression) etc.) are ignored. If the handler is used as
gateway, the ProducerRecord will have its headers enhanced to add the
KafkaHeaders.REPLY_TOPIC unless it already contains such a header. The handler
will not map any additional headers; providing such a payload assumes the headers have
already been mapped.
- Since:
- 5.4
- Author:
- Soby Chacko, Artem Bilan, Gary Russell, Marius Bogoevici, Biju Kunjummen, Tom van den Berge
-
Nested Class Summary
Nested Classes Modifier and Type Class Description static interfaceKafkaProducerMessageHandler.ProducerRecordCreator<K,V>Creates aProducerRecordfrom aMessageand/or properties derived from configuration and/or the message.Nested classes/interfaces inherited from class org.springframework.integration.handler.AbstractReplyProducingMessageHandler
AbstractReplyProducingMessageHandler.RequestHandlerNested classes/interfaces inherited from interface org.springframework.integration.support.management.IntegrationManagement
IntegrationManagement.ManagementOverrides -
Field Summary
Fields inherited from class org.springframework.integration.handler.AbstractMessageProducingHandler
messagingTemplateFields inherited from class org.springframework.integration.context.IntegrationObjectSupport
EXPRESSION_PARSER, loggerFields inherited from interface org.springframework.integration.support.management.IntegrationManagement
METER_PREFIX, RECEIVE_COUNTER_NAME, SEND_TIMER_NAMEFields inherited from interface org.springframework.core.Ordered
HIGHEST_PRECEDENCE, LOWEST_PRECEDENCE -
Constructor Summary
Constructors Constructor Description KafkaProducerMessageHandler(org.springframework.kafka.core.KafkaTemplate<K,V> kafkaTemplate) -
Method Summary
Modifier and Type Method Description protected voiddoInit()StringgetComponentType()Subclasses may implement this method to provide component type information.protected MessageChannelgetFuturesChannel()org.springframework.kafka.support.KafkaHeaderMappergetHeaderMapper()org.springframework.kafka.core.KafkaTemplate<?,?>getKafkaTemplate()protected MessageChannelgetSendFailureChannel()protected MessageChannelgetSendSuccessChannel()protected ObjecthandleRequestMessage(Message<?> message)Subclasses must implement this method to handle the request Message.booleanisRunning()voidprocessSendResult(Message<?> message, org.apache.kafka.clients.producer.ProducerRecord<K,V> producerRecord, ListenableFuture<org.springframework.kafka.support.SendResult<K,V>> future, MessageChannel metadataChannel)voidsetErrorMessageStrategy(ErrorMessageStrategy errorMessageStrategy)Set the error message strategy implementation to use when sending error messages after send failures.voidsetFlushExpression(Expression flushExpression)Specify a SpEL expression that evaluates to aBooleanto determine whether the producer should be flushed after the send.voidsetFuturesChannel(MessageChannel futuresChannel)Set the futures channel.voidsetFuturesChannelName(String futuresChannelName)Set the futures channel name.voidsetHeaderMapper(org.springframework.kafka.support.KafkaHeaderMapper headerMapper)Set the header mapper to use.voidsetMessageKeyExpression(Expression messageKeyExpression)voidsetPartitionIdExpression(Expression partitionIdExpression)voidsetProducerRecordCreator(KafkaProducerMessageHandler.ProducerRecordCreator<K,V> producerRecordCreator)Set aKafkaProducerMessageHandler.ProducerRecordCreatorto create theProducerRecord.voidsetReplyMessageConverter(org.springframework.kafka.support.converter.RecordMessageConverter messageConverter)Set a message converter for gateway replies.voidsetReplyPayloadType(Type payloadType)When using a type-aware message converter (such asStringJsonMessageConverter, set the payload type the converter should create.voidsetSendFailureChannel(MessageChannel sendFailureChannel)Set the failure channel.voidsetSendFailureChannelName(String sendFailureChannelName)Set the failure channel name.voidsetSendSuccessChannel(MessageChannel sendSuccessChannel)Set the success channel.voidsetSendSuccessChannelName(String sendSuccessChannelName)Set the Success channel name.voidsetSendTimeout(long sendTimeout)Specify a timeout in milliseconds for how long thisKafkaProducerMessageHandlershould wait wait for send operation results.voidsetSendTimeoutExpression(Expression sendTimeoutExpression)Specify a SpEL expression to evaluate a timeout in milliseconds for how long thisKafkaProducerMessageHandlershould wait wait for send operation results.voidsetSync(boolean sync)Abooleanindicating if theKafkaProducerMessageHandlershould wait for the send operation results or not.voidsetTimeoutBuffer(int timeoutBuffer)Set a buffer, in milliseconds, added to the configureddelivery.timeout.msto determine the minimum time to wait for the send future completion whensyncis true.voidsetTimestampExpression(Expression timestampExpression)Specify a SpEL expression to evaluate a timestamp that will be added in the Kafka record.voidsetTopicExpression(Expression topicExpression)voidstart()voidstop()Methods inherited from class org.springframework.integration.handler.AbstractReplyProducingMessageHandler
doInvokeAdvisedRequestHandler, getBeanClassLoader, getIntegrationPatternType, getRequiresReply, handleMessageInternal, hasAdviceChain, onInit, setAdviceChain, setBeanClassLoader, setRequiresReplyMethods inherited from class org.springframework.integration.handler.AbstractMessageProducingHandler
addNotPropagatedHeaders, createOutputMessage, getNotPropagatedHeaders, getOutputChannel, isAsync, messageBuilderForReply, produceOutput, resolveErrorChannel, sendErrorMessage, sendOutput, sendOutputs, setAsync, setNotPropagatedHeaders, setOutputChannel, setOutputChannelName, shouldCopyRequestHeaders, shouldSplitOutput, updateNotPropagatedHeadersMethods inherited from class org.springframework.integration.handler.AbstractMessageHandler
handleMessage, onComplete, onError, onNext, onSubscribeMethods inherited from class org.springframework.integration.handler.MessageHandlerSupport
buildSendTimer, destroy, getManagedName, getManagedType, getMetricsCaptor, getOrder, getOverrides, isLoggingEnabled, registerMetricsCaptor, sendTimer, setLoggingEnabled, setManagedName, setManagedType, setOrder, setShouldTrack, shouldTrackMethods inherited from class org.springframework.integration.context.IntegrationObjectSupport
afterPropertiesSet, extractTypeIfPossible, generateId, getApplicationContext, getApplicationContextId, getBeanDescription, getBeanFactory, getBeanName, getChannelResolver, getComponentName, getConversionService, getExpression, getIntegrationProperties, getIntegrationProperty, getMessageBuilderFactory, getTaskScheduler, isInitialized, setApplicationContext, setBeanFactory, setBeanName, setChannelResolver, setComponentName, setConversionService, setMessageBuilderFactory, setPrimaryExpression, setTaskScheduler, toStringMethods inherited from class java.lang.Object
clone, equals, finalize, getClass, hashCode, notify, notifyAll, wait, wait, waitMethods inherited from interface org.springframework.integration.support.management.IntegrationManagement
getThisAsMethods inherited from interface org.springframework.integration.support.context.NamedComponent
getBeanName, getComponentName
-
Constructor Details
-
Method Details
-
setTopicExpression
-
setMessageKeyExpression
-
setPartitionIdExpression
-
setTimestampExpression
Specify a SpEL expression to evaluate a timestamp that will be added in the Kafka record. The resulting value should be aLongtype representing epoch time in milliseconds.- Parameters:
timestampExpression- theExpressionfor timestamp to wait for result fo send operation.- Since:
- 2.3
-
setFlushExpression
Specify a SpEL expression that evaluates to aBooleanto determine whether the producer should be flushed after the send. Defaults to looking for aBooleanvalue in aKafkaIntegrationHeaders.FLUSHheader; false if absent.- Parameters:
flushExpression- theExpression.- Since:
- 3.3
-
setHeaderMapper
public void setHeaderMapper(org.springframework.kafka.support.KafkaHeaderMapper headerMapper)Set the header mapper to use.- Parameters:
headerMapper- the mapper; can be null to disable header mapping.- Since:
- 2.3
-
getHeaderMapper
public org.springframework.kafka.support.KafkaHeaderMapper getHeaderMapper() -
getKafkaTemplate
public org.springframework.kafka.core.KafkaTemplate<?,?> getKafkaTemplate() -
setSync
public void setSync(boolean sync)Abooleanindicating if theKafkaProducerMessageHandlershould wait for the send operation results or not. Defaults tofalse. Insyncmode a downstream send operation exception will be re-thrown.- Parameters:
sync- the send mode; async by default.- Since:
- 2.0.1
-
setSendTimeout
public final void setSendTimeout(long sendTimeout)Specify a timeout in milliseconds for how long thisKafkaProducerMessageHandlershould wait wait for send operation results. Defaults to the kafkadelivery.timeout.msproperty + 5 seconds. The timeout is applied Also applies when sending to the success or failure channels.- Overrides:
setSendTimeoutin classAbstractMessageProducingHandler- Parameters:
sendTimeout- the timeout to wait for result for a send operation.- Since:
- 2.0.1
-
setSendTimeoutExpression
Specify a SpEL expression to evaluate a timeout in milliseconds for how long thisKafkaProducerMessageHandlershould wait wait for send operation results. Defaults to the kafkadelivery.timeout.msproperty + 5 seconds. The timeout is applied only insyncmode. If this expression yields a result that is less than that value, the higher value is used.- Parameters:
sendTimeoutExpression- theExpressionfor timeout to wait for result for a send operation.- Since:
- 2.1.1
- See Also:
setTimeoutBuffer(int)
-
setSendFailureChannel
Set the failure channel. After a send failure, anErrorMessagewill be sent to this channel with a payload of aKafkaSendFailureExceptionwith the failed message and cause.- Parameters:
sendFailureChannel- the failure channel.- Since:
- 2.1.2
-
setSendFailureChannelName
Set the failure channel name. After a send failure, anErrorMessagewill be sent to this channel name with a payload of aKafkaSendFailureExceptionwith the failed message and cause.- Parameters:
sendFailureChannelName- the failure channel name.- Since:
- 2.1.2
-
setSendSuccessChannel
Set the success channel.- Parameters:
sendSuccessChannel- the Success channel.- Since:
- 3.0.2
-
setSendSuccessChannelName
Set the Success channel name.- Parameters:
sendSuccessChannelName- the Success channel name.- Since:
- 3.0.2
-
setFuturesChannel
Set the futures channel.- Parameters:
futuresChannel- the futures channel.- Since:
- 5.4
-
setFuturesChannelName
Set the futures channel name.- Parameters:
futuresChannelName- the futures channel name.- Since:
- 5.4
-
setErrorMessageStrategy
Set the error message strategy implementation to use when sending error messages after send failures. Cannot be null.- Parameters:
errorMessageStrategy- the implementation.- Since:
- 2.1.2
-
setReplyMessageConverter
public void setReplyMessageConverter(org.springframework.kafka.support.converter.RecordMessageConverter messageConverter)Set a message converter for gateway replies.- Parameters:
messageConverter- the converter.- Since:
- 3.0.2
- See Also:
setReplyPayloadType(Type)
-
setReplyPayloadType
When using a type-aware message converter (such asStringJsonMessageConverter, set the payload type the converter should create. Defaults toObject.- Parameters:
payloadType- the type.- Since:
- 3.0.2
- See Also:
setReplyMessageConverter(RecordMessageConverter)
-
setProducerRecordCreator
public void setProducerRecordCreator(KafkaProducerMessageHandler.ProducerRecordCreator<K,V> producerRecordCreator)Set aKafkaProducerMessageHandler.ProducerRecordCreatorto create theProducerRecord.- Parameters:
producerRecordCreator- the creator.- Since:
- 3.2.1
-
setTimeoutBuffer
public void setTimeoutBuffer(int timeoutBuffer)Set a buffer, in milliseconds, added to the configureddelivery.timeout.msto determine the minimum time to wait for the send future completion whensyncis true.- Parameters:
timeoutBuffer- the buffer.- Since:
- 5.4
- See Also:
setSendTimeoutExpression(Expression)
-
getComponentType
Description copied from class:IntegrationObjectSupportSubclasses may implement this method to provide component type information.- Specified by:
getComponentTypein interfaceNamedComponent- Overrides:
getComponentTypein classMessageHandlerSupport
-
getSendFailureChannel
-
getSendSuccessChannel
-
getFuturesChannel
-
doInit
protected void doInit()- Overrides:
doInitin classAbstractReplyProducingMessageHandler
-
start
public void start()- Specified by:
startin interfaceLifecycle- Specified by:
startin interfaceManageableLifecycle
-
stop
public void stop()- Specified by:
stopin interfaceLifecycle- Specified by:
stopin interfaceManageableLifecycle
-
isRunning
public boolean isRunning()- Specified by:
isRunningin interfaceLifecycle- Specified by:
isRunningin interfaceManageableLifecycle
-
handleRequestMessage
Description copied from class:AbstractReplyProducingMessageHandlerSubclasses must implement this method to handle the request Message. The return value may be a Message, a MessageBuilder, or any plain Object. The base class will handle the final creation of a reply Message from any of those starting points. If the return value is null, the Message flow will end here.- Specified by:
handleRequestMessagein classAbstractReplyProducingMessageHandler- Parameters:
message- The request message.- Returns:
- The result of handling the message, or
null.
-
processSendResult
public void processSendResult(Message<?> message, org.apache.kafka.clients.producer.ProducerRecord<K,V> producerRecord, ListenableFuture<org.springframework.kafka.support.SendResult<K,V>> future, MessageChannel metadataChannel) throws InterruptedException, ExecutionException
-