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

使用Kafka DSL的Spring Cloud Stream单元测试遇连接错误及无法终止问题

问题解决:EmbeddedKafka完全支持Kafka DSL,你的问题出在这些细节上

首先明确告诉你:EmbeddedKafka完全支持Kafka DSL(KStreams),你遇到的Connection to node -1 could not be established错误和测试无法终止问题,都是因为几个容易忽略的配置和代码细节导致的,我帮你逐一拆解修复:

1. 绑定接口的注解完全写错了!

你看你的MyBinding接口里,把输出通道OUTPUT也标记成了@Input,这就相当于告诉Spring Cloud Stream这是个输入通道,根本没法发送处理后的消息,还会导致KStreams拓扑构建失败:

// 错误写法
@Input(OUTPUT) 
KStream<Integer, String> messagesOut();

赶紧改成@Output,而且要保证泛型和处理器方法的返回类型一致(你的处理器用的是KStream<Integer,byte[]>,接口也要对应):

// 正确写法
@Output(OUTPUT)
KStream<Integer, byte[]> messagesOut();

2. EmbeddedKafka的配置没传全,KStreams找不到Broker

对于KStreams应用,除了bootstrap-servers和zkNodes,你还需要做两件事:

  • 把Spring Cloud Stream的input/output绑定到你在EmbeddedKafka里创建的so0544in/so0544out topic
  • 给KStreams设置唯一的application.id(这是KStreams应用的必填参数)

在setup()方法里新增这两行系统属性:

// 绑定输入输出通道到EmbeddedKafka的topic
System.setProperty("spring.cloud.stream.bindings.input.destination", "so0544in");
System.setProperty("spring.cloud.stream.bindings.output.destination", "so0544out");
// 给KStreams应用设置唯一ID
System.setProperty("spring.cloud.stream.kafka.streams.binder.applicationId", "test-kstreams-app");

3. 测试无法终止?因为KStreams后台线程没关

KStreams应用会启动后台线程处理流数据,测试结束后如果不主动关闭,进程会一直挂着。你可以在tearDown()方法里添加关闭逻辑:

@Autowired
private StreamsBuilderFactoryBean streamsBuilderFactoryBean;

@After
public void tearDown() {
    if (consumer != null){
        consumer.close();
    }
    // 关闭KStreams客户端,释放线程
    if (streamsBuilderFactoryBean != null && streamsBuilderFactoryBean.isRunning()) {
        streamsBuilderFactoryBean.stop();
    }
}

4. 泛型类型不匹配,序列化会报错

你的处理器方法返回的是KStream<Integer,byte[]>,但绑定接口里定义的是KStream<Integer, String>,类型不匹配会导致消息序列化/反序列化失败,一定要把接口的泛型改成和处理器一致。

修正后的完整可运行代码

绑定接口修正

public interface MyBinding {
    String INPUT = "input";
    String OUTPUT = "output";

    @Input(INPUT)
    KStream<Integer, byte[]> messagesIn();

    @Output(OUTPUT)
    KStream<Integer, byte[]> messagesOut();
}

测试类完整代码

@RunWith(SpringRunner.class)
@SpringBootTest(classes=ExampleKafkaEmbeddedTest.ExampleAppNotWorking.class)
@DirtiesContext(classMode=ClassMode.AFTER_EACH_TEST_METHOD)
public class ExampleKafkaEmbeddedTest {

    @ClassRule
    public static KafkaEmbedded embeddedKafka = new KafkaEmbedded(1, false, "so0544in","so0544out");

    @Autowired
    private KafkaTemplate<Integer, byte[]> template;

    @Autowired
    private KafkaProperties properties;

    @Autowired
    private StreamsBuilderFactoryBean streamsBuilderFactoryBean;

    private static Consumer<String, String> consumer;

    @BeforeClass
    public static void setup() throws Exception{
        System.setProperty("spring.kafka.bootstrap-servers", embeddedKafka.getBrokersAsString());
        System.setProperty("spring.cloud.stream.kafka.binder.zkNodes", embeddedKafka.getZookeeperConnectionString());
        // 新增:绑定通道到EmbeddedKafka的topic
        System.setProperty("spring.cloud.stream.bindings.input.destination", "so0544in");
        System.setProperty("spring.cloud.stream.bindings.output.destination", "so0544out");
        // 新增:KStreams必填的应用ID
        System.setProperty("spring.cloud.stream.kafka.streams.binder.applicationId", "test-kstreams-app");
        System.setProperty("server.port","0");
        System.setProperty("spring.jmx.enabled" , "false");

        Map<String, Object> consumerProps = KafkaTestUtils.consumerProps("group-id", "false", embeddedKafka);
        consumerProps.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
        consumerProps.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
        consumerProps.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");

        DefaultKafkaConsumerFactory<String, String> cf = new DefaultKafkaConsumerFactory<>(consumerProps);
        consumer = cf.createConsumer();
        embeddedKafka.consumeFromAnEmbeddedTopic(consumer, "so0544out");
    }

    @After
    public void tearDown() {
        if (consumer != null){
            consumer.close();
        }
        // 关闭KStreams客户端,解决测试无法终止问题
        if (streamsBuilderFactoryBean != null && streamsBuilderFactoryBean.isRunning()) {
            streamsBuilderFactoryBean.stop();
        }
    }

    @Test
    public void testSendReceive() {
        template.send("so0544in", "foo".getBytes());
        ConsumerRecord<String, String> cr = KafkaTestUtils.getSingleRecord(consumer, "so0544out");
        System.out.println("Contenu chaine resultat : " + cr.value());
        assertEquals("FOO", cr.value());
    }

    @SpringBootApplication
    @EnableBinding(Processor.class)
    public static class ExampleAppWorking {
        public static void main(String[] args) {
            SpringApplication.run(ExampleAppWorking.class, args);
        }

        @StreamListener(Processor.INPUT)
        @SendTo(Processor.OUTPUT)
        public String receive(String in) {
            return in.toUpperCase();
        }
    }

    @SpringBootApplication
    @EnableBinding(MyBinding.class)
    public static class ExampleAppNotWorking {
        public static void main(String[] args) {
            SpringApplication.run(ExampleAppNotWorking.class, args);
        }

        @StreamListener
        @SendTo(MyBinding.OUTPUT)
        public KStream<Integer,byte[]> toUpperCase (@Input(MyBinding.INPUT) KStream<Integer,byte[]> in){
            return in.map((key, val) -> KeyValue.pair(key, new String(val).toUpperCase().getBytes()));
        }
    }
}

额外小建议

  1. 你的Elmhurst.RELEASE版本比较老(对应Spring Boot 2.0.x),建议升级到更高版本的Spring Cloud Stream,新版本对KStreams的测试支持更友好,比如可以用@EmbeddedKafka注解替代KStreamEmbedded规则,配置更简洁。
  2. 如果想做更轻量的单元测试,也可以试试Spring Cloud Stream的TestSupportBinder,它不需要真实的Kafka Broker,只需要在@SpringBootTest里加几个配置就能用:
@SpringBootTest(properties = {
    "spring.cloud.stream.default-binder=test",
    "spring.cloud.stream.bindings.input.destination=so0544in",
    "spring.cloud.stream.bindings.output.destination=so0544out"
})

内容的提问来源于stack exchange,提问作者HL'REB

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.29 09:03:53