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

Spring Cloud Stream使用StreamBridge如何实现生产者发送回调

在使用Kafka binder的Spring Cloud Stream框架中,完全支持配置onSuccess、onFailure类型的发送回调,可以实现和原生producer.send(record, new callback(){...})、返回ListenableFuture的KafkaTemplate.send方法一致的效果,常用实现方案如下:

实现方案

方案1:基于StreamBridge返回值添加单消息粒度回调(最贴近原生KafkaTemplate使用习惯)

StreamBridge的send方法本身返回CompletableFuture<SendResult<byte[], byte[]>>类型结果,和KafkaTemplate返回的ListenableFuture能力完全对齐,你可以直接在返回的Future对象上绑定成功、失败逻辑,不需要额外修改全局配置,开箱即用。

代码示例:

import org.springframework.cloud.stream.function.StreamBridge;
import org.springframework.kafka.support.SendResult;
import org.springframework.stereotype.Component;
import java.util.concurrent.CompletableFuture;

@Component
public class EventPublisher {
    private final StreamBridge streamBridge;

    public EventPublisher(StreamBridge streamBridge) {
        this.streamBridge = streamBridge;
    }

    public void publishEvent(String bindingName, Object eventPayload) {
        CompletableFuture<SendResult<byte[], byte[]>> sendFuture = streamBridge.send(bindingName, eventPayload);
        
        // 绑定发送回调
        sendFuture.whenComplete((result, ex) -> {
            if (ex != null) {
                // onFailure 发送失败逻辑:可执行告警、本地消息表落库、重试等补偿操作
                System.err.println("消息发送失败:" + ex.getMessage());
                return;
            }
            // onSuccess 发送成功逻辑:可执行日志记录、业务状态更新等操作
            System.out.printf("消息发送成功,topic:%s,分区:%d,偏移量:%d%n",
                    result.getRecordMetadata().topic(),
                    result.getRecordMetadata().partition(),
                    result.getRecordMetadata().offset());
        });
    }
}

如果需要同步获取发送结果,直接调用sendFuture.get()即可,行为和KafkaTemplate返回的ListenableFuture的get方法完全一致。

方案2:配置绑定级别的全局发送回调

如果需要对某个Kafka binder下所有生产消息统一处理回调逻辑(比如全局埋点、监控统计),可以通过自定义ProducerListener实现,不需要在每次调用send时单独编写回调。
实现步骤:

  • 自定义类实现org.springframework.kafka.support.ProducerListener接口,重写onSuccess、onError方法编写全局回调逻辑
  • 将自定义Listener注册到对应绑定的生产者工厂中即可生效

代码示例:

import org.apache.kafka.clients.producer.ProducerRecord;
import org.apache.kafka.clients.producer.RecordMetadata;
import org.springframework.kafka.support.ProducerListener;
import org.springframework.stereotype.Component;

@Component
public class GlobalProducerCallback implements ProducerListener<byte[], byte[]> {
    @Override
    public void onSuccess(ProducerRecord<byte[], byte[]> record, RecordMetadata metadata) {
        // 全局发送成功逻辑
    }

    @Override
    public void onError(ProducerRecord<byte[], byte[]> record, RecordMetadata metadata, Exception exception) {
        // 全局发送失败逻辑
    }
}

注意事项

  • 回调触发时机为Kafka Broker确认消息写入完成后,和原生Kafka生产者回调的触发时机完全一致,不存在提前回调的问题
  • 两种方案可以叠加使用,全局回调和单消息回调不会互相冲突

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.26 20:06:26