Spring Boot整合Kafka Streams通过application.properties配置时抛ClassCastException
问题:Spring Boot整合Kafka Streams在K8s部署时的ClassCastException问题
两种Kafka Streams配置方式
方式一:@Bean手动配置(K8s部署运行正常)
通过@Bean显式创建KafkaStreamsConfiguration实例,代码如下:
@Bean(name = KafkaStreamsDefaultConfiguration.DEFAULT_STREAMS_CONFIG_BEAN_NAME) KafkaStreamsConfiguration kStreamsConfig() { Map<String, Object> props = new HashMap<>(); props.put(APPLICATION_ID_CONFIG, "kafka-streams"); props.put(BOOTSTRAP_SERVERS_CONFIG, "localhost:9092"); return new KafkaStreamsConfiguration(props); }
方式二:application.properties自动配置(K8s部署抛出异常)
通过Spring Boot配置文件配置,代码如下:
# application.properties spring.kafka.streams.bootstrap-servers=localhost:9092 spring.kafka.streams.application-id=kafka-streams
异常信息
采用方式二在K8s部署时,应用启动抛出ClassCastException,异常栈如下:
20220728 18:08:17 [XTRA-KAFKA-PRPDUCER11 main] class java.util.ArrayList cannot be cast to class java.lang.String (java.util.ArrayList and java.lang.String are in module java.base of loader 'bootstrap') java.lang.ClassCastException: class java.util.ArrayList cannot be cast to class java.lang.String (java.util.ArrayList and java.lang.String are in module java.base of loader 'bootstrap') java.lang.ClassCastException org.apache.kafka.clients.consumer.KafkaConsumer.<init>(KafkaConsumer.java:666) org.apache.kafka.streams.processor.internals.DefaultKafkaClientSupplier.getRestoreConsumer(DefaultKafkaClientSupplier.java:49) org.apache.kafka.streams.processor.internals.StreamThread.create(StreamThread.java:343) org.apache.kafka.streams.KafkaStreams.createAndAddStreamThread(KafkaStreams.java:956) org.apache.kafka.streams.KafkaStreams.<init>(KafkaStreams.java:948) org.apache.kafka.streams.KafkaStreams.<init>(KafkaStreams.java:845) org.apache.kafka.streams.KafkaStreams.<init>(KafkaStreams.java:751) org.springframework.kafka.config.StreamsBuilderFactoryBean.start(StreamsBuilderFactoryBean.java:349) org.springframework.context.support.DefaultLifecycleProcessor.doStart(DefaultLifecycleProcessor.java:178) org.springframework.context.support.DefaultLifecycleProcessor.access$200(DefaultLifecycleProcessor.java:54) org.springframework.context.support.DefaultLifecycleProcessor$LifecycleGroup.start(DefaultLifecycleProcessor.java:356) java.base/java.lang.Iterable.forEach(Iterable.java:75) org.springframework.context.support.DefaultLifecycleProcessor.startBeans(DefaultLifecycleProcessor.java:155) org.springframework.context.support.DefaultLifecycleProcessor.onRefresh(DefaultLifecycleProcessor.java:123) org.springframework.context.support.AbstractApplicationContext.finishRefresh(AbstractApplicationContext.java:935) org.springframework.context.support.AbstractApplicationContext.refresh(AbstractApplicationContext.java:586) org.springframework.boot.web.servlet.context.ServletWebServerApplicationContext.refresh(ServletWebServerApplicationContext.java:147) org.springframework.boot.SpringApplication.refresh(SpringApplication.java:734) org.springframework.boot.SpringApplication.refreshContext(SpringApplication.java:408) org.springframework.boot.SpringApplication.run(SpringApplication.java:308) org.springframework.boot.SpringApplication.run(SpringApplication.java:1306) org.springframework.boot.SpringApplication.run(SpringApplication.java:1295) com.example.springkafka.SpringKafkaApplication.main(SpringKafkaApplication.java:10) java.base/jdk.internal.reflect.DirectMethodHandleAccessor.invoke(DirectMethodHandleAccessor.java:104) java.base/java.lang.reflect.Method.invoke(Method.java:577) org.springframework.boot.loader.MainMethodRunner.run(MainMethodRunner.java:49) org.springframework.boot.loader.Launcher.launch(Launcher.java:108) org.springframework.boot.loader.Launcher.launch(Launcher.java:58) org.springframework.boot.loader.JarLauncher.main(JarLauncher.java:65)
注:上述
bootstrap-servers为示例,实际应用中并非localhost:9092
场景补充
本地IntelliJ开发环境中,两种配置方式均无异常,仅在K8s部署时出现该问题。
原因分析
- 核心问题是K8s环境中,
spring.kafka.streams.bootstrap-servers被解析为ArrayList类型,而Kafka客户端期望接收String类型(逗号分隔的地址列表)。 - 手动@Bean配置时,直接传入的是String类型参数,不存在类型转换问题;而通过application.properties自动配置时,若K8s的配置源(如ConfigMap、环境变量)将bootstrap-servers以数组/列表形式传递,Spring Boot自动绑定配置时会错误地将集合类型赋值给期望String的属性,最终导致Kafka消费者初始化时抛出类型转换异常。
解决方法
- 调整K8s配置传递格式:确保
spring.kafka.streams.bootstrap-servers的值为逗号分隔的字符串,而非数组格式。例如在ConfigMap中配置为spring.kafka.streams.bootstrap-servers=kafka-0:9092,kafka-1:9092。 - 自定义配置绑定逻辑:创建自定义配置类,手动将集合类型的bootstrap-servers转换为逗号分隔的字符串后注入到Kafka Streams配置中。
- 保留@Bean手动配置方式:继续使用显式@Bean配置,完全控制参数类型,避免自动配置的类型绑定问题。
内容的提问来源于stack exchange,提问作者Danny.an
相关产品推荐
相关产品推荐

