使用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/so0544outtopic - 给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())); } } }
额外小建议
- 你的Elmhurst.RELEASE版本比较老(对应Spring Boot 2.0.x),建议升级到更高版本的Spring Cloud Stream,新版本对KStreams的测试支持更友好,比如可以用
@EmbeddedKafka注解替代KStreamEmbedded规则,配置更简洁。 - 如果想做更轻量的单元测试,也可以试试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
相关产品推荐
相关产品推荐

