升级Spring Boot Starter Parent至3.2.x及以上版本后Embedded Kafka启动失败
升级Spring Boot Starter Parent至3.2.x及以上版本后Embedded Kafka启动失败
我正在尝试使用org.springframework.kafka.test.context.EmbeddedKafka为Kafka消费者运行集成测试,目前让spring-boot-starter-parent负责依赖版本管理,以下是pom.xml文件:
<parent> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-parent</artifactId> <version>3.2.4</version> <relativePath/> <!-- lookup parent from repository --> </parent> <properties> <java.version>17</java.version> </properties> <dependencies> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-web</artifactId> </dependency> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-test</artifactId> <scope>test</scope> </dependency> <dependency> <groupId>org.springframework.kafka</groupId> <artifactId>spring-kafka</artifactId> </dependency> <dependency> <groupId>org.springframework.kafka</groupId> <artifactId>spring-kafka-test</artifactId> <scope>test</scope> </dependency> </dependencies>
Kafka消费者代码
@Configuration public class KafkaEventListener { @RetryableTopic( attempts = "#{'${kafka.max.retry.attempts}'}", autoCreateTopics = "#{'${kafka.auto.create.retry.topics}'}", backoff = @Backoff( delayExpression = "#{'${kafka.retry.init-interval}'}", multiplierExpression = "#{'${kafka.retry.backoff.multiplier}'}"), include = { Exception.class }, timeout = "#{'${kafka.max.retry.duration}'}", topicSuffixingStrategy = TopicSuffixingStrategy.SUFFIX_WITH_INDEX_VALUE, dltStrategy = DltStrategy.FAIL_ON_ERROR) @KafkaListener(topics = "${kafka.topic.test}", groupId = "${kafka.group.id.test}", containerFactory = "testKafkaListenerContainerFactory") public void listen(@Payload MetadataMessage input, @Header(KafkaHeaders.OFFSET) String offset) { System.out.println(input.getValue()); } @DltHandler public void deadLetterHandler(@Payload(required = false) MetadataMessage data, @Header(KafkaHeaders.RECEIVED_TOPIC) String topic) { System.out.println(String.format("Event from topic %s has been dead-lettered. Event data : %s", topic, data.toString())); } }
Kafka配置类
@Configuration public class KafkaConsumerConfig { @Value(value = "${spring.kafka.bootstrap-servers}") private String bootstrapAddress; @Value(value = "${kafka.group.id.test}") private String groupId; private Map<String, Object> consumerFactoryConfigs() { Map<String, Object> props = new HashMap<>(); props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapAddress); props.put(ConsumerConfig.GROUP_ID_CONFIG, groupId); props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, JsonDeserializer.class); // Disabled kafka auto acknowledgement to gain more flexibility props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false); props.put(ConsumerConfig.MAX_PARTITION_FETCH_BYTES_CONFIG, "20971520"); props.put(ConsumerConfig.FETCH_MAX_BYTES_CONFIG, "20971520"); props.put(JsonDeserializer.TRUSTED_PACKAGES, "*"); return props; } @Bean public ConsumerFactory<String, MetadataMessage> metadataConsumerFactory() { return new DefaultKafkaConsumerFactory<>(consumerFactoryConfigs(), new StringDeserializer(), new ErrorHandlingDeserializer<>(new JsonDeserializer<>(MetadataMessage.class))); } @Bean public ConcurrentKafkaListenerContainerFactory<String, MetadataMessage> testKafkaListenerContainerFactory() { ConcurrentKafkaListenerContainerFactory<String, MetadataMessage> factory = new ConcurrentKafkaListenerContainerFactory<>(); factory.setConsumerFactory(metadataConsumerFactory()); factory.getContainerProperties().setAckMode(ContainerProperties.AckMode.MANUAL_IMMEDIATE); return factory; } }
测试类
@SpringBootTest(classes = EmbeddedKafkaApplication.class, webEnvironment = SpringBootTest.WebEnvironment.RANDOM_PORT) @TestInstance(TestInstance.Lifecycle.PER_CLASS) @TestPropertySource(locations = { "classpath:application.properties" }) @EmbeddedKafka(partitions = 1, brokerProperties = { "listeners=PLAINTEXT://localhost:29092", "port=29092" }) public class EmbeddedKafkaTest { @Autowired private KafkaEventListener kafkaEventListener; @Test public void testKafkaEvent() { kafkaEventListener.listen(new MetadataMessage("kafka message from test"), "1", mock(Acknowledgment.class)); } }
当使用spring-boot-starter-parent 3.1.10版本时,测试可以正常运行。但切换到3.2.0或更高版本(比如3.2.4)后,测试就失败了。
3.1.10版本启动时的日志(部分)
. ____ _ __ _ _ /\\ / ___'_ __ _ _(_)_ __ __ _ \ \ \ \ ( ( )\___ | '_ | '_| | '_ \/ _` | \ \ \ \ \\/ ___)| |_)| | | | | || (_| | ) ) ) ) ' |____| .__|_| |_|_| |_\__, | / / / / =========|_|==============|___/=/_/_/_/ :: Spring Boot :: (v3.1.10) 2024-03-31T16:37:37.354+05:30 INFO 26280 --- [ main] k.utils.Log4jControllerRegistration$ : Registered kafka:type=kafka.Log4jController MBean 2024-03-31T16:37:37.448+05:30 INFO 26280 --- [ main] o.a.zookeeper.server.ZooKeeperServer : 2024-03-31T16:37:37.448+05:30 INFO 26280 --- [ main] o.a.zookeeper.server.ZooKeeperServer : ______ _ 2024-03-31T16:37:37.448+05:30 INFO 26280 --- [ main] o.a.zookeeper.server.ZooKeeperServer : |___ / | | 2024-03-31T16:37:37.448+05:30 INFO 26280 --- [ main] o.a.zookeeper.server.ZooKeeperServer : / / ___ ___ | | __ ___ ___ _ __ ___ _ __ 2024-03-31T16:37:37.448+05:30 INFO 26280 --- [ main] o.a.zookeeper.server.ZooKeeperServer : / / / _ \ / _ \ | |/ / / _ \ / _ \ | '_ \ / _ \ | '__| 2024-03-31T16:37:37.448+05:30 INFO 26280 --- [ main] o.a.zookeeper.server.ZooKeeperServer : / /__ | (_) | | (_) | | < | __/ | __/ | |_) | | __/ | | 2024-03-31T16:37:37.448+05:30 INFO 26280 --- [ main] o.a.zookeeper.server.ZooKeeperServer : /_____| \___/ \___/ |_|\_\ \___| \___| | .__/ \___| |_| 2024-03-31T16:37:37.448+05:30 INFO 26280 --- [ main] o.a.zookeeper.server.ZooKeeperServer : | | 2024-03-31T16:37:37.448+05:30 INFO 26280 --- [ main] o.a.zookeeper.server.ZooKeeperServer : |_| 2024-03-31T16:37:37.448+05:30 INFO 26280 --- [ main] o.a.zookeeper.server.ZooKeeperServer : 2024-03-31T16:37:37.460+05:30 INFO 26280 --- [ main] o.a.zookeeper.server.ZooKeeperServer : Server environment:zookeeper.version=3.6.4--d65253dcf68e9097c6e95a126463fd5fdeb4521c, built on 12/18/2022 18:10 GMT 2024-03-31T16:37:37.460+05:30 INFO 26280 --- [ main] o.a.zookeeper.server.ZooKeeperServer : Server environment:host.name=SADEEP-M.Zone24x7.lk 2024-03-31T16:37:37.460+05:30 INFO 26280 --- [ main] o.a.zookeeper.server.ZooKeeperServer : Server environment:java.version=17.0.8.1 2024-03-31T16:37:37.460+05:30 INFO 26280 --- [ main] o.a.zookeeper.server.ZooKeeperServer : Server environment:java.vendor=Amazon.com Inc.
3.1.10版本测试成功执行的最后日志
2024-03-31T16:37:41.551+05:30 INFO 26280 --- [ner#0-dlt-0-C-1] o.s.k.l.KafkaMessageListenerContainer : testGroupId-dlt: partitions assigned: [testTopic-dlt-0] 2024-03-31T16:37:41.551+05:30 INFO 26280 --- [0-retry-1-0-C-1] o.s.k.l.KafkaMessageListenerContainer : testGroupId-retry-1: partitions assigned: [testTopic-retry-1-0] 2024-03-31T16:37:41.551+05:30 INFO 26280 --- [ntainer#0-0-C-1] o.s.k.l.KafkaMessageListenerContainer : testGroupId: partitions assigned: [testTopic-0] 2024-03-31T16:37:41.551+05:30 INFO 26280 --- [0-retry-0-0-C-1] o.s.k.l.KafkaMessageListenerContainer : testGroupId-retry-0: partitions assigned: [testTopic-retry-0-0] kafka message from test
可以清楚看到ZooKeeperServer在运行。
但切换到3.2.x版本后,看不到ZooKeeper相关日志,测试也无法执行,以下是3.2.4版本的部分日志:
2024-03-31T16:46:29.072+05:30 INFO 25024 --- [embedded-kafka] [ main] o.a.kafka.common.utils.AppInfoParser : Kafka version: 3.6.1 2024-03-31T16:46:29.072+05:30 INFO 25024 --- [embedded-kafka] [ main] o.a.kafka.common.utils.AppInfoParser : Kafka commitId: 5e3c2b738d253ff5 2024-03-31T16:46:29.072+05:30 INFO 25024 --- [embedded-kafka] [ main] o.a.kafka.common.utils.AppInfoParser : Kafka startTimeMs: 1711883789072 2024-03-31T16:46:29.075+05:30 INFO 25024 --- [embedded-kafka] [| adminclient-2] org.apache.kafka.clients.NetworkClient : [AdminClient clientId=adminclient-2] Node -1 disconnected. 2024-03-31T16:46:29.075+05:30 WARN 25024 --- [embedded-kafka] [| adminclient-2] org.apache.kafka.clients.NetworkClient : [AdminClient clientId=adminclient-2] Connection to node -1 (localhost/127.0.0.1:29092) could not be established. Broker may not be available. 2024-03-31T16:46:29.191+05:30 INFO 25024 --- [embedded-kafka] [| adminclient-2] org.apache.kafka.clients.NetworkClient : [AdminClient clientId=adminclient-2] Node -1 disconnected. 2024-03-31T16:46:29.191+05:30 WARN 25024 --- [embedded-kafka] [| adminclient-2] org.apache.kafka.clients.NetworkClient : [AdminClient clientId=adminclient-2] Connection to node -1 (localhost/127.0.0.1:29092) could not be established. Broker may not be available. 2024-03-31T16:46:29.300+05:30 INFO 25024 --- [embedded-kafka] [| adminclient-2] org.apache.kafka.clients.NetworkClient : [AdminClient clientId=adminclient-2] Node -1 disconnected. 2024-03-31T16:46:29.300+05:30 WARN 25024 --- [embedded-kafka] [| adminclient-2] org.apache.kafka.clients.NetworkClient : [AdminClient clientId=adminclient-2] Connection to node -1 (localhost/127.0.0.1:29092) could not be established. Broker may not be available. 2024-03-31T16:46:29.518+05:30 INFO 25024 --- [embedded-kafka] [| adminclient-2] org.apache.kafka.clients.NetworkClient : [AdminClient clientId=adminclient-2] Node -1 disconnected.
另外启动时也看不到ZooKeeper的标志性日志:
/\\ / ___'_ __ _ _(_)_ __ __ _ \ \ \ \ ( ( )\___ | '_ | '_| | '_ \/ _` | \ \ \ \ \\/ ___)| |_)| | | | | || (_| | ) ) ) ) ' |____| .__|_| |_|_| |_\__, | / / / / =========|_|==============|___/=/_/_/_/ :: Spring Boot :: (v3.2.4) 2024-03-31T16:46:26.369+05:30 INFO 25024 --- [embedded-kafka] [ main] k.utils.Log4jControllerRegistration$ : Registered kafka:type=kafka.Log4jController MBean 2024-03-31T16:46:26.387+05:30 INFO 25024 --- [embedded-kafka] [ main] org.apache.zookeeper.common.X509Util : Setting -D jdk.tls.rejectClientInitiatedRenegotiation=true to disable client-init
相关产品推荐
相关产品推荐

