【发布时间】:2020-12-25 21:08:24
【问题描述】:
我已经将spring集成kafka用于一个简单的read-process-write场景。消息驱动端点 bean 的配置如下:
<kafka:message-driven-channel-adapter listener-container="listenerContainer"
channel="processChannel"
mode="batch"/>
<bean id="listenerContainer" class="org.springframework.kafka.listener.KafkaMessageListenerContainer" parent="kafkaMessageListenerContainerAbstract">
<constructor-arg>
<bean class="org.springframework.kafka.listener.ContainerProperties">
<constructor-arg name="topics" value="test"/>
<property name="transactionManager" ref="kafkaTransactionManager"/>
<property name="eosMode" value="BETA"/>
</bean>
</constructor-arg>
</bean>
当 KafkaMessageListenerContainer 轮询一批记录(例如 10 条记录)时,将启动一个 kafka 事务,但 IntegrationBatchMessageListener 会将所有 10 条记录整合为一条消息强>
message = toMessagingMessage(records, acknowledgment, consumer);
是否有任何解决方案可以在单个 kafka 事务中处理批处理,但在消息驱动的端点中单独处理每条记录,然后提交事务?
如果可能的话,我不想使用手动确认。
我见过 BatchToRecordAdapter 类,但认为它在消息驱动端点中不可用。
【问题讨论】:
标签: apache-kafka spring-integration spring-kafka