如何在Spring Integration Kafka中手动确认消费者读取的消息
Alright, let's figure out how to implement manual offset commits for your Spring Integration Kafka setup (3.1.2.RELEASE) where you want full control over when offsets are committed—specifically after successful transformation, and no commit if transformation fails. Here's a step-by-step solution with code examples tailored to your configuration:
1. Update Listener Container Configuration
First, you need to configure your KafkaMessageListenerContainer to use manual acknowledgment mode. This tells Spring Kafka to hold off on auto-committing offsets and wait for your explicit signal.
Modify your ContainerProperties bean to add the ackMode property:
<bean id="container1" class="org.springframework.kafka.listener.KafkaMessageListenerContainer"> <constructor-arg> <bean class="org.springframework.kafka.core.DefaultKafkaConsumerFactory"> <constructor-arg> <map> <entry key="bootstrap.servers" value="${spring.kafka.bootstrap-servers}" /> <entry key="enable.auto.commit" value="false" /> <entry key="auto.commit.interval.ms" value="100" /> <entry key="session.timeout.ms" value="15000" /> <entry key="group.id" value="${spring.kafka.consumer.group-id}" /> <entry key="key.deserializer" value="org.apache.kafka.common.serialization.StringDeserializer" /> <entry key="value.deserializer" value="com.test.CustomDeserializer" /> </map> </constructor-arg> </bean> </constructor-arg> <constructor-arg> <bean class="org.springframework.kafka.listener.ContainerProperties"> <constructor-arg name="topics" value="${spring.kafka.topics}" /> <!-- Enable manual acknowledgment mode --> <property name="ackMode" value="MANUAL" /> <!-- Optional: Add retry template if you want to retry failed messages before sending to error channel --> <!-- <property name="retryTemplate" ref="retryTemplate" /> --> </bean> </constructor-arg> </bean>
2. Access Acknowledgment in Your Transformer
When using manual acknowledgment mode, Spring Integration automatically passes the Kafka Acknowledgment object as a message header (KafkaHeaders.ACKNOWLEDGMENT). Your transformer can retrieve this header to commit offsets only after successful processing.
You have two clean ways to access the acknowledgment:
Option 1: Receive the full Message<?> object
@Component("transformerBean") public class CustomTransformer { public Object transform(Message<?> message) { // Get your deserialized custom object from the message payload CustomObject payload = (CustomObject) message.getPayload(); // Retrieve the acknowledgment header Acknowledgment acknowledgment = message.getHeaders() .get(KafkaHeaders.ACKNOWLEDGMENT, Acknowledgment.class); try { // Execute your transformation logic TransformedObject transformedResult = performTransformation(payload); // Commit offset only if transformation succeeds if (acknowledgment != null) { acknowledgment.acknowledge(); } return transformedResult; } catch (Exception e) { // Transformation failed: don't commit, throw exception to route to error channel throw new MessageTransformationException(message, "Failed to transform message", e); } } private TransformedObject performTransformation(CustomObject payload) { // Your custom transformation logic here } }
Option 2: Use @Header annotation to inject the acknowledgment directly
@Component("transformerBean") public class CustomTransformer { public TransformedObject transform( CustomObject payload, @Header(KafkaHeaders.ACKNOWLEDGMENT) Acknowledgment acknowledgment) { try { // Execute your transformation logic TransformedObject transformedResult = performTransformation(payload); // Commit offset on success acknowledgment.acknowledge(); return transformedResult; } catch (Exception e) { // Transformation failed: don't commit, throw exception to trigger error handling throw new RuntimeException("Transformation failed for payload: " + payload.getId(), e); } } private TransformedObject performTransformation(CustomObject payload) { // Your custom transformation logic here } }
3. Error Handling Notes
- If your transformer throws an exception, the message will be routed to your configured
error-channel="errorChannel". Since you didn't callacknowledge(), the offset won't be committed, and Spring Kafka will re-deliver the message (based on your container's retry settings). - You can use your existing
recoveryCallbackanderror-message-strategyto handle retries or dead-lettering for failed messages without committing offsets.
Key Points to Remember
ackMode="MANUAL"tells the container to wait for your explicitacknowledge()call before committing offsets. UseMANUAL_IMMEDIATEinstead if you want offsets to be committed immediately after callingacknowledge()(instead of waiting for the next batch interval).- Ensure your
message-converter(if custom) preserves Kafka headers—defaultMessagingMessageConverterdoes this automatically, so you don't need extra configuration if you're using it.
内容的提问来源于stack exchange,提问作者Muthu

