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
A Message Handler for Apache Kafka; 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.
The handler also 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 ClassesModifier and TypeClassDescriptionstatic interfaceCreates 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
ConstructorsConstructorDescriptionKafkaProducerMessageHandler(org.springframework.kafka.core.KafkaTemplate<K, V> kafkaTemplate) -
Method Summary
Modifier and TypeMethodDescriptionprotected voiddoInit()Subclasses may implement this method to provide component type information.protected MessageChannelorg.springframework.kafka.support.KafkaHeaderMapperorg.springframework.kafka.core.KafkaTemplate<?,?> protected MessageChannelprotected MessageChannelprotected ObjecthandleRequestMessage(Message<?> message) Subclasses must implement this method to handle the request Message.booleanvoidprocessSendResult(Message<?> message, org.apache.kafka.clients.producer.ProducerRecord<K, V> producerRecord, CompletableFuture<org.springframework.kafka.support.SendResult<K, V>> future, MessageChannel metadataChannel) voidsetAssignmentDuration(Duration assignmentDuration) Set the time to wait for partition assignment, when used as a gateway, to determine the default reply-to topic/partition.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.final voidsetSendTimeout(long sendTimeout) Specify a timeout in milliseconds for how long thisKafkaProducerMessageHandlershould wait for send operation results.voidsetSendTimeoutExpression(Expression sendTimeoutExpression) Specify a SpEL expression to evaluate a timeout in milliseconds for how long thisKafkaProducerMessageHandlershould 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) voidsetUseTemplateConverter(boolean useTemplateConverter) Set to true to use the template's message converter to create theProducerRecordinstead of theproducerRecordCreator.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, setupMessageProcessor, shouldCopyRequestHeaders, shouldSplitOutput, updateNotPropagatedHeadersMethods inherited from class org.springframework.integration.handler.AbstractMessageHandler
handleMessage, onComplete, onError, onNext, onSubscribe, setObservationConventionMethods inherited from class org.springframework.integration.handler.MessageHandlerSupport
buildSendTimer, destroy, getManagedName, getManagedType, getMetricsCaptor, getObservationRegistry, getOrder, getOverrides, isLoggingEnabled, isObserved, registerMetricsCaptor, registerObservationRegistry, 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 reactor.core.CoreSubscriber
currentContextMethods inherited from interface org.springframework.integration.support.management.IntegrationManagement
getThisAsMethods inherited from interface org.springframework.integration.support.context.NamedComponent
getBeanName, getComponentName
-
Constructor Details
-
KafkaProducerMessageHandler
-
-
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.
-
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.
-
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.
-
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.
-
setSendTimeout
public final void setSendTimeout(long sendTimeout) Specify a timeout in milliseconds for how long thisKafkaProducerMessageHandlershould 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.
-
setSendTimeoutExpression
Specify a SpEL expression to evaluate a timeout in milliseconds for how long thisKafkaProducerMessageHandlershould 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.- See Also:
-
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.
-
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.
-
setSendSuccessChannel
Set the success channel.- Parameters:
sendSuccessChannel- the Success channel.
-
setSendSuccessChannelName
Set the Success channel name.- Parameters:
sendSuccessChannelName- the Success channel name.
-
setFuturesChannel
Set the futures channel.- Parameters:
futuresChannel- the futures channel.
-
setFuturesChannelName
Set the futures channel name.- Parameters:
futuresChannelName- the futures channel name.
-
setErrorMessageStrategy
Set the error message strategy implementation to use when sending error messages after send failures. Cannot be null.- Parameters:
errorMessageStrategy- the implementation.
-
setReplyMessageConverter
public void setReplyMessageConverter(org.springframework.kafka.support.converter.RecordMessageConverter messageConverter) Set a message converter for gateway replies.- Parameters:
messageConverter- the converter.- See Also:
-
setReplyPayloadType
When using a type-aware message converter (such asStringJsonMessageConverter, set the payload type the converter should create. Defaults toObject.- Parameters:
payloadType- the type.- See Also:
-
setProducerRecordCreator
public void setProducerRecordCreator(KafkaProducerMessageHandler.ProducerRecordCreator<K, V> producerRecordCreator) Set aKafkaProducerMessageHandler.ProducerRecordCreatorto create theProducerRecord. Ignored ifuseTemplateConverteris true.- Parameters:
producerRecordCreator- the creator.- See Also:
-
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.- See Also:
-
setUseTemplateConverter
public void setUseTemplateConverter(boolean useTemplateConverter) Set to true to use the template's message converter to create theProducerRecordinstead of theproducerRecordCreator.- Parameters:
useTemplateConverter- true to use the converter.- Since:
- 5.5.5
- See Also:
-
setAssignmentDuration
Set the time to wait for partition assignment, when used as a gateway, to determine the default reply-to topic/partition.- Parameters:
assignmentDuration- the assignmentDuration to set.- Since:
- 6.0
-
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, CompletableFuture<org.springframework.kafka.support.SendResult<K, throws InterruptedException, ExecutionExceptionV>> future, MessageChannel metadataChannel)
-