Embedded Kafka JUnit测试执行失败求助
问题
本地运行Embedded Kafka的JUnit测试时执行失败,相关代码如下:
@RunWith(SpringRunner.class) @SpringBootTest @DirtiesContext @EmbeddedKafka(partitions = 1, topics = {EmbeddedKafkaIntegrationTest.TEST_TOPIC},bootstrapServersProperty = "spring.kafka.bootstrap-servers") @ActiveProfiles("test") public class EmbeddedKafkaIntegrationTest { static final String TEST_TOPIC = "MERCHANT-SERVICE-BILLABLE-EVENTS-DPG-IN"; @Value("${spring.kafka.consumer.group-id}") private String groupId; @Value("${spring.kafka.consumer.auto-offset-reset}") private String offsetReset; @Autowired private MessagePublisher producer; @Autowired private KafkaConsumer consumer; @Autowired private EmbeddedKafkaBroker embeddedKafka; @Captor ArgumentCaptor<ConsumerRecord<String, BillableEventsRequest>> billableEventsRequestArgumentCaptor; @Captor ArgumentCaptor<String> topicArgumentCaptor; static { System.setProperty(EmbeddedKafkaBroker.BROKER_LIST_PROPERTY, "spring.kafka.bootstrap-servers"); } @Test public void embeddedKafka_whenSendingToProducer_thenMessageReceived() throws IOException, ParseException { //Producer BillableEventsRequest event = createTestRequest(); producer.publish(event); //consumer verify(consumer, timeout(1000).times(1)).consume(billableEventsRequestArgumentCaptor.capture()); ConsumerRecord<String, BillableEventsRequest> payload = billableEventsRequestArgumentCaptor.getValue(); assertNotNull(payload); assertTrue(TEST_TOPIC.contains(topicArgumentCaptor.getValue())); testEvents(payload,event); } }
使用环境版本
- Spring Boot : v2.3.5.RELEASE
- Kafka version: 2.5.1
- Apache Camel 3.9.0
- Java : v11.0.15
- zookeeper.version=3.5.9-83df9301aa5c2a5d284a9940177808c01bc35cef, built on 01/06/2021 20:03 GMT
测试错误日志
- 配置
spring.kafka.producer.bootstrap-servers为IP地址:9092时:
[Consumer clientId=consumer-sbilling-billable-events-dpg-in-consumer-1, groupId=sbilling-billable-events-dpg-in-consumer] Connection to node -1 (/10.165.101.110:9092) could not be established. Broker may not be available. 2022-10-17 17:49:50.643 WARN 1820 --- [ntainer#0-0-C-1] org.apache.kafka.clients.NetworkClient : [Consumer clientId=consumer-sbilling-billable-events-dpg-in-consumer-1, groupId=sbilling-billable-events-dpg-in-consumer] Bootstrap broker 10.165.101.110:9092 (id: -1 rack: null) disconnected 2022-10-17 17:49:52.006 WARN 1820 --- [ad | producer-1] org.apache.kafka.clients.NetworkClient : [Producer clientId=producer-1] Connection to node -1 (INPNQLT896845.in.db.com/10.165.101.110:9092) could not be established. Broker may not be available. 2022-10-17 17:49:52.006 WARN 1820 --- [ad | producer-1] org.apache.kafka.clients.NetworkClient : [Producer clientId=producer-1] Bootstrap broker 10.165.101.110:9092 (id: -1 rack: null) disconnected
- 配置
spring.kafka.producer.bootstrap-servers为localhost:9092时:
2022-10-17 17:19:26.743 WARN 27272 --- [ntainer#0-0-C-1] org.apache.kafka.clients.NetworkClient : [Consumer clientId=consumer-sbilling-billable-events-dpg-in-consumer-1, groupId=sbilling-billable-events-dpg-in-consumer] Connection to node -1 (localhost/127.0.0.1:9092) could not be established. Broker may not be available. 2022-10-17 17:19:26.743 WARN 27272 --- [ntainer#0-0-C-1] org.apache.kafka.clients.NetworkClient : [Consumer clientId=consumer-sbilling-billable-events-dpg-in-consumer-1, groupId=sbilling-billable-events-dpg-in-consumer] Bootstrap broker localhost:9092 (id: -1 rack: null) disconnected 2022-10-17 17:19:27.687 WARN 27272 --- [ad | producer-1] org.apache.kafka.clients.NetworkClient : [Producer clientId=producer-1] Connection to node -1 (localhost/127.0.0.1:9092) could not be established. Broker may not be available. 2022-10-17 17:19:27.688 WARN 27272 --- [ad | producer-1] org.apache.kafka.clients.NetworkClient : [Producer clientId=producer-1] Bootstrap broker localhost:9092 (id: -1 rack: null) disconnected
- 最终测试报错:
[ERROR] Tests run: 1, Failures: 0, Errors: 1, Skipped: 0, Time elapsed: 73.128 s <<< FAILURE! - in com.db.spricing.ms.broker.EmbeddedKafkaIntegrationTest [ERROR] embeddedKafka_whenSendingToProducer_thenMessageReceived Time elapsed: 60.24 s <<< ERROR! org.springframework.kafka.KafkaException: Send failed; nested exception is org.apache.kafka.common.errors.TimeoutException: Topic test.in.topic not present in metadata after 60000 ms. at com.db.spricing.ms.broker.EmbeddedKafkaIntegrationTest.embeddedKafka_whenSendingToProducer_thenMessageReceived(EmbeddedKafkaIntegrationTest.java:89) Caused by: org.apache.kafka.common.errors.TimeoutException: Topic test.in.topic not present in metadata after 60000 ms.
排查与解决方案
核心问题分析
生产者、消费者均无法连接Embedded Kafka,且最终报错的test.in.topic与代码定义的TEST_TOPIC不匹配,说明存在配置注入失效和Topic配置不一致两个核心问题。
具体修复步骤
移除冗余的静态配置代码块
@EmbeddedKafka注解中已经通过bootstrapServersProperty = "spring.kafka.bootstrap-servers"自动将Embedded Kafka的随机端口注入到配置项,无需再通过System.setProperty重复设置,直接删除以下代码:static { System.setProperty(EmbeddedKafkaBroker.BROKER_LIST_PROPERTY, "spring.kafka.bootstrap-servers"); }统一Topic配置
报错中出现的test.in.topic与代码定义的MERCHANT-SERVICE-BILLABLE-EVENTS-DPG-IN不一致,需:- 检查
MessagePublisher实现类,确认其发送消息使用的Topic是TEST_TOPIC常量; - 查看测试环境配置文件(
application-test.yml/application-test.properties),确保生产者的Topic配置与测试类中定义的一致。
- 检查
删除硬编码的Broker地址
如果测试配置文件中手动设置了spring.kafka.bootstrap-servers=localhost:9092或固定IP,直接删除该配置。@EmbeddedKafka会自动生成随机可用端口并注入到该属性,硬编码会导致程序连接外部不存在的Kafka,而非Embedded实例。调整测试超时时间
测试中verify(consumer, timeout(1000).times(1))的1秒超时过短,Embedded Kafka初始化和消息传递需要更多时间,建议调整为5秒:verify(consumer, timeout(5000).times(1)).consume(billableEventsRequestArgumentCaptor.capture());校验Camel Kafka组件配置
由于使用了Apache Camel,需确认Camel的Kafka组件是否正确读取spring.kafka.bootstrap-servers配置,避免组件内部硬编码Broker地址。
内容的提问来源于stack exchange,提问作者Trups

