Kafka Streams内部重分区主题在Nomad环境未创建问题求助
Kafka Streams内部重分区主题创建失败排查问题
问题场景
- 运行一个Kafka Streams应用,拓扑中两个分支的窗口聚合操作会生成2个内部重分区主题
- 本地或docker-compose环境运行正常,消费者读取时AdminClient会自动创建这些内部主题
- 部署到Nomad编排环境(Kafka Broker启用SSL)后,尽管全局流配置已提供SSL证书,应用仍无法创建内部重分区主题
已确认信息
- 内部AdminClient配置正常,日志输出的配置如下:
bootstrap.servers = [kafka-01.aws:9093, kafka-02.aws:9093, kafka-03.aws:9093] client.dns.lookup = use_all_dns_ips client.id = namespace.application-cd79b0a0-ffe2-4694-a70e-d023d7d12908-admin connections.max.idle.ms = 300000 default.api.timeout.ms = 60000 metadata.max.age.ms = 300000 metric.reporters = [] metrics.num.samples = 2 metrics.recording.level = INFO metrics.sample.window.ms = 30000 receive.buffer.bytes = 65536 reconnect.backoff.max.ms = 1000 reconnect.backoff.ms = 50 request.timeout.ms = 30000 retries = 2147483647 retry.backoff.ms = 100 sasl.client.callback.handler.class = null sasl.jaas.config = null sasl.kerberos.kinit.cmd = /usr/bin/kinit sasl.kerberos.min.time.before.relogin = 60000 sasl.kerberos.service.name = null sasl.kerberos.ticket.renew.jitter = 0.05 sasl.kerberos.ticket.renew.window.factor = 0.8 sasl.login.callback.handler.class = null sasl.login.class = null sasl.login.connect.timeout.ms = null sasl.login.read.timeout.ms = null sasl.login.refresh.buffer.seconds = 300 sasl.login.refresh.min.period.seconds = 60 sasl.login.refresh.window.factor = 0.8 sasl.login.refresh.window.jitter = 0.05 sasl.login.retry.backoff.max.ms = 10000 sasl.login.retry.backoff.ms = 100 sasl.mechanism = GSSAPI sasl.oauthbearer.clock.skew.seconds = 30 sasl.oauthbearer.expected.audience = null sasl.oauthbearer.expected.issuer = null sasl.oauthbearer.jwks.endpoint.refresh.ms = 3600000 sasl.oauthbearer.jwks.endpoint.retry.backoff.max.ms = 10000 sasl.oauthbearer.jwks.endpoint.retry.backoff.ms = 100 sasl.oauthbearer.jwks.endpoint.url = null sasl.oauthbearer.scope.claim.name = scope sasl.oauthbearer.sub.claim.name = sub sasl.oauthbearer.token.endpoint.url = null security.protocol = SSL security.providers = null send.buffer.bytes = 131072 socket.connection.setup.timeout.max.ms = 30000 socket.connection.setup.timeout.ms = 10000 ssl.cipher.suites = null ssl.enabled.protocols = [TLSv1.2, TLSv1.3] ssl.endpoint.identification.algorithm = ssl.engine.factory.class = null ssl.key.password = null ssl.keymanager.algorithm = SunX509 ssl.keystore.certificate.chain = null ssl.keystore.key = null ssl.keystore.location = /tmp/kafka_keystore.jks ssl.keystore.password = [hidden] ssl.keystore.type = JKS ssl.protocol = TLSv1.3 ssl.provider = null ssl.secure.random.implementation = null ssl.trustmanager.algorithm = PKIX ssl.truststore.certificates = null ssl.truststore.location = /tmp/kafka_truststore.jks ssl.truststore.password = [hidden] ssl.truststore.type = JKS
- 应用终止前的错误日志:无AdminClient创建内部主题的记录,仅输出主题不存在的警告:
[namespace.application-cd79b0a0-ffe2-4694-a70e-d023d7d12908-StreamThread-1] WARN org.apache.kafka.clients.NetworkClient - [Consumer clientId=namespace.application-cd79b0a0-ffe2-4694-a70e-d023d7d12908-StreamThread-1-consumer, groupId=namespace.application] Error while fetching metadata with correlation id 8 : {namespace.application-GroupedResults_2_Store-repartition=UNKNOWN_TOPIC_OR_PARTITION, namespace.application-GroupedResults_1_Store-repartition=UNKNOWN_TOPIC_OR_PARTITION} [namespace.application-cd79b0a0-ffe2-4694-a70e-d023d7d12908-StreamThread-1] WARN org.apache.kafka.clients.NetworkClient - [Consumer clientId=namespace.application-cd79b0a0-ffe2-4694-a70e-d023d7d12908-StreamThread-1-consumer, groupId=namespace.application] Error while fetching metadata with correlation id 9 : {namespace.application-GroupedResults_2_Store-repartition=UNKNOWN_TOPIC_OR_PARTITION, namespace.application-GroupedResults_1_Store-repartition=UNKNOWN_TOPIC_OR_PARTITION} [namespace.application-cd79b0a0-ffe2-4694-a70e-d023d7d12908-StreamThread-1] WARN org.apache.kafka.clients.NetworkClient - [Consumer clientId=namespace.application-cd79b0a0-ffe2-4694-a70e-d023d7d12908-StreamThread-1-consumer, groupId=namespace.application] Error while fetching metadata with correlation id 10 : {namespace.application-GroupedResults_2_Store-repartition=UNKNOWN_TOPIC_OR_PARTITION, namespace.application-GroupedResults_1_Store-repartition=UNKNOWN_TOPIC_OR_PARTITION} [namespace.application-cd79b0a0-ffe2-4694-a70e-d023d7d12908-StreamThread-1] WARN org.apache.kafka.clients.NetworkClient - [Consumer clientId=namespace.application-cd79b0a0-ffe2-4694-a70e-d023d7d12908-StreamThread-1-consumer, groupId=namespace.application] Error while fetching metadata with correlation id 11 : {namespace.application-GroupedResults_2_Store-repartition=UNKNOWN_TOPIC_OR_PARTITION, namespace.application-.GroupedResults_1_Store-repartition=UNKNOWN_TOPIC_OR_PARTITION}
- 自定义AdminClient正常工作:SpringBoot中自定义的AdminClient Bean能成功创建输出主题和死信主题,配置与Streams内部AdminClient类似
排查思路
- 权限差异验证:对比自定义AdminClient和Streams内部AdminClient的身份,检查Broker ACL是否允许Streams使用的
client.id对应的身份创建主题。重点验证是否允许对namespace.application-*-repartition这类前缀的主题执行Create操作 - SSL配置细节核对:
- 检查
ssl.endpoint.identification.algorithm为空是否符合Broker要求,若Broker强制SNI验证,会导致连接失败 - 确认Nomad容器内
/tmp/kafka_keystore.jks和/tmp/kafka_truststore.jks的存在性及应用进程的读权限,自定义AdminClient可能使用了不同路径或已提前加载证书 - 验证Broker支持的TLS版本,若Broker仅支持TLSv1.2,而Streams配置
ssl.protocol=TLSv1.3会导致握手失败
- 检查
- 主题创建时机与超时检查:Streams默认在拓扑初始化时创建内部主题,检查Nomad环境下网络延迟是否过高,
default.api.timeout.ms(当前60s)是否足够,或Broker端是否存在请求延迟 - 内部主题命名合法性检查:最后一条日志中出现
namespace.application-.GroupedResults_1_Store-repartition(多了一个点),排查拓扑中分组操作的名称是否包含特殊字符,导致生成的主题名不符合Broker命名规则,创建请求被拒绝但未在应用日志中输出 - Broker日志分析:直接查看Kafka Broker日志,搜索内部主题名的创建请求,确认是否有权限拒绝、命名非法、SSL握手失败等具体错误
- Streams配置覆盖检查:验证SpringBoot自动配置是否覆盖了Kafka Streams的AdminClient相关配置,比如
streams.admin.*前缀的配置是否正确传递到内部AdminClient - Nomad网络策略验证:检查Nomad网络策略是否允许应用容器访问Kafka Broker的9093端口,同时确认
client.dns.lookup=use_all_dns_ips配置在Nomad环境下是否导致DNS解析异常,部分Broker节点无法访问
内容的提问来源于stack exchange,提问作者Bruno Vilhena
相关产品推荐
相关产品推荐

