Spring Cloud Stream生产者autoStartup=false配置不生效问题
Spring Cloud Stream 生产者自动启动禁用失效、Bean延迟初始化问题
问题背景
使用 Spring Cloud Stream 2021.0.2 搭配 Spring Boot 2.6.7 时,遇到生产者自动启动无法禁用的问题;消费者端的同配置经actuator/bindings端点校验可正常生效。
当前使用的配置如下:
spring: kafka: bootstrap-servers: ${KAFKA_URL:http://localhost:29092} cloud: stream: function: definition: mqtt;mqttRoot;kafka bindings: mqtt-out-0: destination: input producer: autoStartup: false mqttRoot-out-0: destination: input producer: autoStartup: false kafka-in-0: destination: output group: mqtt consumer: autoStartup: false concurrency: 1
业务场景
基于spring-cloud-stream-kafka-binder与Hive MQ Java Reactor客户端实现MQTT与Kafka的双向消息桥接。目前已通过ApplicationRunner实现应用启动完成后再连接MQTT的逻辑,需要确认是否有其他方式可延迟生产者/消费者Bean的初始化。
相关实现代码如下:
InputConfig 配置类
@Configuration @Log4j2 public class InputConfig { @Value("${mqtt.subscriptionTopic}") private String subscriptionTopic; private final Mqtt3ReactorClient inboundMqttClient; private final Mqtt3ReactorClient inboundRootMqttClient; @Autowired public InputConfig(@Qualifier("inboundMqttClient") Mqtt3ReactorClient inboundMqttClient, @Qualifier("inboundRootMqttClient") Mqtt3ReactorClient inboundRootMqttClient) { this.inboundMqttClient = inboundMqttClient; this.inboundRootMqttClient = inboundRootMqttClient; } @Bean public Supplier<Flux<byte[]>> mqtt() { return new Input(inboundMqttClient, subscriptionTopic); } @Bean public Supplier<Flux<byte[]>> mqttRoot() { return new Input(inboundRootMqttClient, subscriptionTopic); } }
Input 消息生产实现类
@Log4j2 public class Input implements Supplier<Flux<byte[]>>, ApplicationRunner { private final Mqtt3ReactorClient inboundMqttClient; private final String subscriptionTopic; public Input(Mqtt3ReactorClient inboundMqttClient, String subscriptionTopic) { this.inboundMqttClient = inboundMqttClient; this.subscriptionTopic = subscriptionTopic; } @Override public Flux<byte[]> get() { FluxWithSingle<Mqtt3Publish, Mqtt3SubAck> subAckAndMatchingPublishes = inboundMqttClient.subscribePublishesWith() .topicFilter(subscriptionTopic).qos(MqttQos.AT_LEAST_ONCE) .applySubscribe(); return subAckAndMatchingPublishes .doOnSingle(subAck -> log.info("Subscribed, " + subAck.getReturnCodes())) .doOnNext(publish -> log.debug( "Received publish" + ", topic: " + publish.getTopic() + ", QoS: " + publish.getQos() + ", payload: " + new String(publish.getPayloadAsBytes()))) .map(Mqtt3Publish::getPayloadAsBytes); } @Override public void run(ApplicationArguments args) { Mono<Mqtt3ConnAck> connAckSingle = inboundMqttClient.connect(); connAckSingle .doOnSuccess(connAck -> log.info("Connected, " + connAck.getReturnCode())) .doOnError(throwable -> log.info("Connection failed, " + throwable.getMessage())) .subscribe(); } }
OutputConfig 配置类
@Configuration @Log4j2 public class OutputConfig { @Value("${mqtt.publishTopic}") private String publishTopic; private final Mqtt3ReactorClient outboundMqttClient; private final Mqtt3ReactorClient outboundRootMqttClient; @Autowired public OutputConfig(@Qualifier("outboundMqttClient") Mqtt3ReactorClient outboundMqttClient, @Qualifier("outboundRootMqttClient") Mqtt3ReactorClient outboundRootMqttClient) { this.outboundMqttClient = outboundMqttClient; this.outboundRootMqttClient = outboundRootMqttClient; } @Bean public Consumer<Flux<Output.GatewayNotification>> kafka() { return new Output(outboundMqttClient, outboundRootMqttClient, publishTopic); } }
Output 消息消费实现类
@Log4j2 public class Output implements Consumer<Flux<Output.GatewayNotification>>, ApplicationRunner { private final Mqtt3ReactorClient outboundMqttClient; private final Mqtt3ReactorClient outboundRootMqttClient; private final String publishTopic; public Output(Mqtt3ReactorClient outboundMqttClient, Mqtt3ReactorClient outboundRootMqttClient, String publishTopic) { this.outboundMqttClient = outboundMqttClient; this.outboundRootMqttClient = outboundRootMqttClient; this.publishTopic = publishTopic; } @Override public void accept(Flux<Output.GatewayNotification> gatewayNotifications) { Flux<Mqtt3Publish> messagesToPublish = gatewayNotifications .map(gatewayNotification -> Mqtt3Publish.builder() .topic(publishTopic.replace("\\{gateway_address\\}", gatewayNotification.getAddress())) .qos(MqttQos.AT_LEAST_ONCE) .payload(gatewayNotification.getPayload().getBytes()) .build()); outboundMqttClient.publish(messagesToPublish) .doOnNext(publishResult -> log.debug( "Publish acknowledged: " + new String(publishResult.getPublish().getPayloadAsBytes()))) .subscribe(); outboundRootMqttClient.publish(messagesToPublish) .doOnNext(publishResult -> log.debug( "Publish acknowledged: " + new String(publishResult.getPublish().getPayloadAsBytes()))) .subscribe(); } @Override public void run(ApplicationArguments args) { Mono<Mqtt3ConnAck> connAckSingle = outboundMqttClient.connect(); connAckSingle .doOnSuccess(connAck -> log.info("Connected, " + connAck.getReturnCode())) .doOnError(throwable -> log.info("Connection failed, " + throwable.getMessage())) .subscribe(); Mono<Mqtt3ConnAck> connAckSingleRoot = outboundRootMqttClient.connect(); connAckSingleRoot .doOnSuccess(connAck -> log.info("Connected, " + connAck.getReturnCode())) .doOnError(throwable -> log.info("Connection failed, " + throwable.getMessage())) .subscribe(); } @Data public static class GatewayNotification { private String address; private String payload; private Long buildingId; } }
异常信息
当前触发的底层异常截图如下:
内容的提问来源于stack exchange,提问作者Rod McCutcheon
相关产品推荐
相关产品推荐

