Spring Cloud Stream Kafka Binder下/bindings Actuator停止生产者失效问题
问题原因说明
- 该现象是你使用的Spring Cloud Stream版本(对应Spring Boot 2.5.x的为3.1.x版本)的已知行为,不属于配置错误:
你用到的轮询式Supplier生产者的调度逻辑由Spring Cloud Function内置的TaskScheduler独立管理,而/actuator/bindings端点的STOP操作仅控制对应binding的消息投递能力,二者在该版本中没有做状态联动。
也就是说你把producer-out-0设为STOPPED之后,调度器仍然会按照你配置的5秒间隔调用你的producer方法,所以你会看到持续打印生产日志;如果该版本中停止状态的binding没有做消息拦截,就会出现消息仍然发送到Kafka的现象。 - 为什么PAUSE操作报错:
PAUSE状态仅对消费者binding生效,生产者binding本身不支持暂停操作,只有STOP/RESUME两种状态可用,所以设置PAUSE会返回错误日志。
解决方案
- 方案1:升级依赖版本
升级Spring Boot到2.6.x及以上版本,对应Spring Cloud Stream 3.2.x及以上版本,该版本已经实现了binding状态和Supplier轮询调度的联动:当生产者binding被设为STOPPED时,对应的轮询任务会自动暂停,不会再调用Supplier方法,恢复为STARTED状态时自动恢复调度。 - 方案2:低版本兼容改造
如果你必须保留当前2.5.3的Spring Boot版本,可以自行在Supplier中增加状态开关控制:- 定义一个全局的AtomicBoolean变量标记生产者是否启用
- 可以自定义Actuator端点或者监听binding的状态变更事件来修改这个开关的状态
- 在Supplier方法执行时先判断开关状态,如果是关闭状态直接返回null,Spring Cloud Stream会自动忽略null返回值,不会执行后续的消息投递逻辑,你也可以同步关闭生产日志的打印。
补充验证点
你可以在停止binding后,对比Supplier的生产日志和消费者的消费日志的数量,如果消费者收到的消息数量少于生产日志的数量,说明停止状态的binding已经在拦截消息,只是调度器还在运行而已。
内容的提问来源于stack exchange,提问作者Tomboyo
相关产品推荐
相关产品推荐

