Kafka Streams容器部署报错:Sink节点未识别致数据处理不全
Kafka Streams容器部署后部分数据未处理,报错"sink node is not recognized"
错误日志
18:58:32.647 [kafka-producer-network-thread | my-app-events-processor.splitPackets-cf462b02-f1e3-4ed5-a1e7-acc1f040495b-StreamThread-1-producer] ERROR o.a.k.s.p.i.RecordCollectorImpl - stream-thread [my-app-events-processor.splitPackets-cf462b02-f1e3-4ed5-a1e7-acc1f040495b-StreamThread-1] task [0_0] Unable to records bytes produced to topic my-app.packet.surface by sink node split-server-log as the node is not recognized. Known sink nodes are []. 18:58:49.216 [kafka-producer-network-thread | my-app-events-processor.splitPackets-cf462b02-f1e3-4ed5-a1e7-acc1f040495b-StreamThread-1-producer] ERROR o.a.k.s.p.i.RecordCollectorImpl - stream-thread [my-app-events-processor.splitPackets-cf462b02-f1e3-4ed5-a1e7-acc1f040495b-StreamThread-1] task [0_0] Unable to records bytes produced to topic my-app.packet.surface by sink node split-server-log as the node is not recognized. Known sink nodes are []. 18:59:05.981 [kafka-producer-network-thread | my-app-events-processor.splitPackets-cf462b02-f1e3-4ed5-a1e7-acc1f040495b-StreamThread-1-producer] ERROR o.a.k.s.p.i.RecordCollectorImpl - stream-thread [my-app-events-processor.splitPackets-cf462b02-f1e3-4ed5-a1e7-acc1f040495b-StreamThread-1] task [0_0] Unable to records bytes produced to topic my-app.packet.surface by sink node split-server-log as the node is not recognized. Known sink nodes are []. 19:00:28.484 [my-app-events-processor.splitPackets-cf462b02-f1e3-4ed5-a1e7-acc1f040495b-StreamThread-1] INFO o.a.k.s.p.internals.StreamThread - stream-thread [my-app-events-processor.splitPackets-cf462b02-f1e3-4ed5-a1e7-acc1f040495b-StreamThread-1] Processed 3 total records, ran 0 punctuators, and committed 3 total tasks since the last update
现象描述
本地运行正常,部署到Docker容器后仅部分数据被处理:部分KafkaStream实例可正常生产数据,其余仅能消费。应用预期消费JSON数据并生成供Leaflet Web地图使用的图片,但仅部分实例能完成此流程。
应用配置
- 语言:Kotlin 1.7.20
- 框架:Kafka Streams 3.3.1
- 部署:Docker容器,通过独立Kotlin协程启动4个独立KafkaStream实例
- 依赖环境:同一Docker网络中的Kafka Kraft容器
- 系统环境:Debian GNU/Linux 11 (bullseye),Kernel 5.10.0-18-amd64,x86-64架构
- Docker版本:20.10.19,docker-compose 1.29.2
Kafka配置
Streams实例消费者配置(示例)
18:38:25.138 [DefaultDispatcher-worker-5 @my-app-events-processor.splitPackets#5] INFO o.a.k.s.p.internals.StreamThread - stream-thread [my-app-events-processor.splitPackets-d7b897b3-3a10-48d6-95c7-e291cb1839d8-StreamThread-1] Creating restore consumer client 18:38:25.142 [DefaultDispatcher-worker-5 @my-app-events-processor.splitPackets#5] INFO o.a.k.c.consumer.ConsumerConfig - ConsumerConfig values: allow.auto.create.topics = true auto.commit.interval.ms = 5000 auto.offset.reset = none bootstrap.servers = [http://kafka:29092] check.crcs = true client.dns.lookup = use_all_dns_ips client.id = my-app-events-processor.splitPackets-d7b897b3-3a10-48d6-95c7-e291cb1839d8-StreamThread-1-restore-consumer client.rack = connections.max.idle.ms = 540000 default.api.timeout.ms = 60000 enable.auto.commit = false exclude.internal.topics = true fetch.max.bytes = 52428800 fetch.max.wait.ms = 500 fetch.min.bytes = 1 group.id = null group.instance.id = null heartbeat.interval.ms = 3000 interceptor.classes = [] internal.leave.group.on.close = false internal.throw.on.fetch.stable.offset.unsupported = true isolation.level = read_committed key.deserializer = class org.apache.kafka.common.serialization.ByteArrayDeserializer max.partition.fetch.bytes = 1048576 max.poll.interval.ms = 300000 max.poll.records = 1000 metadata.max.age.ms = 300000 metric.reporters = [] metrics.num.samples = 2 metrics.recording.level = INFO metrics.sample.window.ms = 30000 partition.assignment.strategy = [class org.apache.kafka.clients.consumer.RangeAssignor, class org.apache.kafka.clients.consumer.CooperativeStickyAssignor] receive.buffer.bytes = 65536 reconnect.backoff.max.ms = 1000 reconnect.backoff.ms = 50 request.timeout.ms = 30000 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 = PLAINTEXT security.providers = null send.buffer.bytes = 131072 session.timeout.ms = 45000 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 = https 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 = null ssl.keystore.password = null 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 = null ssl.truststore.password = null ssl.truststore.type = JKS value.deserializer = class org.apache.kafka.common.serialization.ByteArrayDeserializer
Kafka Kraft服务端配置
基于官方模板,仅修改advertised.listeners:
advertised.listeners=PLAINTEXT://kafka:29092,PLAINTEXT_HOST://localhost:9092
Docker配置(docker-compose)
version: "3.9" services: events-processors: image: events-processors container_name: events-processors restart: unless-stopped environment: KAFKA_BOOTSTRAP_SERVERS: "http://kafka:29092" networks: - my-app-infra-nw depends_on: - infra-kafka secrets: - source: my-app_config target: /.secret.config.yml infra-kafka: image: kafka-kraft container_name: infra-kafka restart: unless-stopped networks: my-app-infra-nw: aliases: [ kafka ] volumes: - "./config/kafka-server.properties:/kafka/server.properties" ports: # note: other Docker containers should use 29092 - "9092:9092" - "9093:9093"
问题解答
错误含义
该错误表示Kafka Streams的RecordCollector组件无法找到名为split-server-log的sink节点,无法将数据生产到目标主题my-app.packet.surface,当前实例已知的sink节点列表为空。本质是Streams实例的拓扑定义未正确加载sink节点信息,导致生产流程中断。
解决步骤
1. 统一多实例的拓扑定义
- 确保4个KafkaStreams实例使用完全一致的拓扑代码,检查是否存在条件分支导致部分实例遗漏了
sink node split-server-log的定义。 - 验证协程启动逻辑,保证每个实例初始化拓扑的过程完全相同,没有因并发初始化导致的拓扑差异。
2. 修正Kafka Bootstrap地址
- 配置中
bootstrap.servers使用了http://kafka:29092,但Kafka的PLAINTEXT协议不需要http://前缀,这会导致客户端无法正确解析地址。将其改为kafka:29092或PLAINTEXT://kafka:29092。 - 更新
events-processors容器的环境变量KAFKA_BOOTSTRAP_SERVERS为修正后的地址。
3. 确保主题创建完成后再启动Streams实例
- 通过Kafka Admin创建主题后,添加等待逻辑:轮询主题状态,直到主题存在且所有分区的副本都处于可用状态(
Isr列表包含所有副本)。 - 避免在主题未完全就绪时启动Streams实例,否则实例无法获取完整的主题元数据,导致sink节点无法被识别。
4. 检查Streams实例的application.id唯一性
- 每个KafkaStreams实例必须使用唯一的
application.id,如果多个实例共用同一个ID,会导致消费者组冲突,任务分配异常,表现为部分实例无法正常生产数据。 - 确认4个实例的
application.id各不相同,比如在ID后添加唯一标识(如实例序号、UUID)。
5. 验证Docker网络连通性与Kafka元数据
- 在
events-processors容器内执行nc -zv kafka 29092,测试是否能正常连接到Kafka服务端口。 - 查看Kafka Kraft的日志,确认主题
my-app.packet.surface的分区、副本状态正常,元数据同步无异常。
内容的提问来源于stack exchange,提问作者aSemy
相关产品推荐
相关产品推荐

