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

升级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
相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.17 12:48:48