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
相关产品推荐
相关产品推荐

