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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.28 22:33:39