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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.16 00:41:38