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

Kafka Streams内部重分区主题在Nomad环境未创建问题求助

Kafka Streams内部重分区主题创建失败排查问题

问题场景

  • 运行一个Kafka Streams应用,拓扑中两个分支的窗口聚合操作会生成2个内部重分区主题
  • 本地或docker-compose环境运行正常,消费者读取时AdminClient会自动创建这些内部主题
  • 部署到Nomad编排环境(Kafka Broker启用SSL)后,尽管全局流配置已提供SSL证书,应用仍无法创建内部重分区主题

已确认信息

  1. 内部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
  1. 应用终止前的错误日志:无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}
  1. 自定义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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.28 15:55:09