Spring Cloud Dataflow的kafka-source-kafka连接配置异常排查
问题场景
使用Spring Cloud Dataflow(SCDF)的kafka-source-kafka组件构建流应用,已配置外部Kafka的消费者参数如下:
app.kafka-source-kafka.kafka.supplier.topics=TEST_TOPIC app.kafka-source-kafka.spring.kafka.consumer.bootstrap-servers=xxxxxxxx app.kafka-source-kafka.spring.kafka.consumer.enable-auto-commit=true app.kafka-source-kafka.spring.kafka.consumer.group-id=B2B_EH_PASTHROUGH_FALLBACK app.kafka-source-kafka.spring.kafka.consumer.properties.sasl.jaas.config=org.apache.kafka.common.security.plain.PlainLoginModule required username=\"xxxxxx\" password=\"xxxxxx\"; app.kafka-source-kafka.spring.kafka.consumer.properties.sasl.mechanism=PLAIN app.kafka-source-kafka.spring.kafka.consumer.properties.security.protocol=SASL_SSL app.kafka-source-kafka.spring.kafka.consumer.properties.ssl.enabled.protocols=TLSv1.2 app.kafka-source-kafka.spring.kafka.consumer.properties.ssl.truststore.location=/etc/secrets/kafka.client.truststore.jks app.kafka-source-kafka.spring.kafka.consumer.properties.ssl.truststore.password=xxxx app.kafka-source-kafka.spring.kafka.consumer.properties.ssl.truststore.type=JKS app.kafka-source-kafka.spring.kafka.listener.async-acks=true
应用启动后日志显示已成功加入外部Kafka消费组,但几分钟后消费者的bootstrap-servers自动切换为内部Kafka地址(10.100.172.208:9092),并出现认证失败错误:
Connection to node -1 (my-release-kafka.springdata.svc.cluster.local/10.100.172.208:9092) terminated during authentication. This may happen due to any of the following reasons: (1) Authentication failed
成功启动日志
2024-01-04T05:45:28.226Z INFO 1 --- [ main] o.a.kafka.common.utils.AppInfoParser : Kafka version: 3.4.1 2024-01-04T05:45:28.226Z INFO 1 --- [ main] o.a.kafka.common.utils.AppInfoParser : Kafka commitId: 8a516edc2755df89 2024-01-04T05:45:28.226Z INFO 1 --- [ main] o.a.kafka.common.utils.AppInfoParser : Kafka startTimeMs: 1704347128225 2024-01-04T05:45:28.318Z INFO 1 --- [ main] o.a.k.clients.consumer.KafkaConsumer : [Consumer clientId=consumer-B2B_EH_PASTHROUGH_FALLBACK-1, groupId=B2B_EH_PASTHROUGH_FALLBACK] Subscribed to topic(s): EIS.TOPIC.PASS.EH.IN.DIT 2024-01-04T05:45:28.409Z INFO 1 --- [ main] s.i.k.i.KafkaMessageDrivenChannelAdapter : started bean 'kafkaMessageDrivenChannelAdapterSpec'; defined in: 'class path resource [org/springframework/cloud/fn/supplier/kafka/KafkaSupplierConfiguration.class]'; from source: 'org.springframework.cloud.fn.supplier.kafka.KafkaSupplierConfiguration.kafkaMessageDrivenChannelAdapterSpec(org.springframework.cloud.fn.supplier.kafka.KafkaSupplierProperties,org.springframework.kafka.config.ConcurrentKafkaListenerContainerFactory,org.springframework.beans.factory.ObjectProvider,org.springframework.beans.factory.ObjectProvider,org.springframework.beans.factory.ObjectProvider,org.springframework.cloud.fn.common.config.ComponentCustomizer)' 2024-01-04T05:45:28.521Z INFO 1 --- [ main] o.s.b.w.embedded.tomcat.TomcatWebServer : Tomcat started on port(s): 8080 (http) with context path '' 2024-01-04T05:45:28.715Z INFO 1 --- [ main] .c.s.a.k.s.k.KafkaSourceKafkaApplication : Started KafkaSourceKafkaApplication in 31.902 seconds (process running for 34.939)
排查方向
外部Kafka的advertised.listeners配置错误
Kafka客户端初始连接bootstrap地址成功后,会从Broker获取集群元数据,其中包含advertised.listeners配置的地址,后续客户端会切换到这个地址进行通信。如果外部Kafka的advertised.listeners配置的是内部集群地址(比如K8S内部Service地址),而你的应用没有配置该地址的认证信息,就会触发认证失败。需要检查外部Kafka Broker的advertised.listeners配置,确保其对外暴露的是可访问的外部地址。SCDF全局配置覆盖应用级参数
检查SCDF部署环境是否默认安装了内部Kafka,且存在全局配置(比如spring.cloud.stream.kafka.binder.bootstrap-servers)优先级高于你设置的应用级参数。SCDF的全局应用配置、环境变量可能会覆盖你单独为kafka-source-kafka设置的参数,导致后续连接切换到内部Kafka。Kafka Source配置前缀错误
kafka-source-kafka是基于Spring Cloud Function的Supplier组件,其配置前缀应为kafka.supplier而非spring.kafka.consumer。你当前配置的是app.kafka-source-kafka.spring.kafka.consumer.bootstrap-servers,可能未生效,或者被组件自身的kafka.supplier.bootstrap-servers配置覆盖。尝试替换为app.kafka-source-kafka.kafka.supplier.bootstrap-servers并重新部署。认证会话重建时的地址获取异常
初始认证成功后,会话过期需要重建连接时,Broker返回了内部地址。这种情况可能是Broker的listener配置存在内外网区分,而客户端在重建时未正确获取外部地址。可以检查Broker的listener配置是否区分了内外网,并确保外部客户端能获取到正确的外部listener地址。网络/DNS解析不稳定
外部Kafka的bootstrap地址在运行一段时间后无法解析,导致客户端 fallback到内部Kafka地址。可以检查DNS解析是否稳定,或者直接使用IP地址配置外部Kafka的bootstrap-servers进行测试。
内容的提问来源于stack exchange,提问作者Zishnu

