You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

使用Kafka出站通道发送消息后如何获取偏移量

获取Kafka输出通道适配器发送后的偏移量

因为你已经配置了sync="true"(同步发送模式),Spring Integration的Kafka输出通道适配器会等待消息发送完成,并将Kafka的SendResult对象作为处理结果返回,你可以通过以下两种方式获取偏移量:

1. 通过ExpressionEvaluatingRequestHandlerAdvice的成功表达式获取

你的配置中已经使用了ExpressionEvaluatingRequestHandlerAdvice,可以直接在onSuccessExpression中通过#result变量访问SendResult对象,从中提取偏移量、主题、分区等元数据:

比如修改你的requestHandlerAdvice配置,调整onSuccessExpression:

<bean id="requestHandlerAdvice"
      class="org.springframework.integration.handler.advice.ExpressionEvaluatingRequestHandlerAdvice">
    <property name="trapException" value="true"/>
    <!-- 直接在表达式中获取并打印偏移量 -->
    <property name="onSuccessExpression" value="T(System).out.println('消息发送成功:主题=' + #result.recordMetadata.topic() + ', 分区=' + #result.recordMetadata.partition() + ', 偏移量=' + #result.recordMetadata.offset())"/>
    <property name="successChannelName" value="successChannel"/>
    <property name="onFailureExpression" ref="failure"/>
    <property name="failureChannelName" value="failureChannel"/>
</bean>

如果是使用已有的success表达式bean,只需确保表达式中能访问#result变量即可,示例:

<bean id="success" class="org.springframework.expression.common.LiteralExpression">
    <constructor-arg value="'发送成功,偏移量:' + #result.recordMetadata.offset()"/>
</bean>

2. 通过成功通道(successChannel)获取

发送成功后,适配器会将结果发送到successChannel,此时消息的IntegrationMessageHeaderAccessor.RESULT_HEADER头中存储了SendResult对象,你可以编写消息处理器提取偏移量:

Java代码示例:

import org.springframework.integration.support.IntegrationMessageHeaderAccessor;
import org.springframework.kafka.support.SendResult;
import org.springframework.messaging.Message;
import org.springframework.integration.annotation.ServiceActivator;

@ServiceActivator(inputChannel = "successChannel")
public void handleKafkaSendSuccess(Message<?> message) {
    // 从消息头中获取SendResult对象
    SendResult<?, ?> sendResult = (SendResult<?, ?>) message.getHeaders()
            .get(IntegrationMessageHeaderAccessor.RESULT_HEADER);
    
    if (sendResult != null) {
        long offset = sendResult.getRecordMetadata().offset();
        String topic = sendResult.getRecordMetadata().topic();
        int partition = sendResult.getRecordMetadata().partition();
        
        // 这里可根据业务需求处理偏移量,比如入库、日志记录等
        System.out.printf("消息发送完成:主题=%s,分区=%d,偏移量=%d%n", topic, partition, offset);
    }
}

注意事项

  • 仅当sync="true"时,适配器才会返回SendResult并传递到advice或success通道;若使用异步模式(sync="false"),则需要通过KafkaTemplate的自定义回调处理结果。
  • 同步模式下会等待Kafka的ack响应,会增加发送延迟,适合需要确认发送结果的场景。

内容的提问来源于stack exchange,提问作者TechUk

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.19 19:50:17