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

升级Apache Flink Python Kafka Connector至1.16+引发客户端断开问题

问题描述

我基于Confluent cp-kafka 7.2部署Apache Kafka Broker,搭配ZooKeeper、Logstash消费者,以及使用librdkafka的自定义Zeek生产者。当使用Apache Flink 1.16+(任意高于1.15.4的版本)及对应flink-sql-connector-kafka Jar包在Python中创建消费者时,出现以下异常:

  • Logstash消费者与Zeek生产者会断开连接,且数分钟内无法恢复
  • Broker完全无法响应CLI查询(如kafka-consumer-groups命令)

切换至Python 3.7搭配Flink 1.15.4时,该断开问题不会出现。我已尝试Python 3.7、3.9、3.10搭配Flink 1.15.4至1.18版本,仅Python 3.7+Flink 1.15.4组合正常,新版本均触发问题。

Kafka Broker日志

---logstash related---
[2024-04-01 23:58:18,586] INFO [GroupCoordinator 1]: Assignment received from leader logstash-0-4bb31f6d-f4e0-4196-afb9-c3be6d251a77 for group logstash for generation 20. The group has 1 members, 0 of which are static. (kafka.coordinator.group.GroupCoordinator)
[2024-04-01 23:58:33,598] INFO [GroupCoordinator 1]: Member logstash-0-4bb31f6d-f4e0-4196-afb9-c3be6d251a77 in group logstash has failed, removing it from the group (kafka.coordinator.group.GroupCoordinator)
---zeek related---
[2024-04-01 23:58:54,164] ERROR [KafkaApi-1] Unexpected error handling request RequestHeader(apiKey=INIT_PRODUCER_ID, apiVersion=4, clientId=producer-16, correlationId=2) -- InitProducerIdRequestData(transactionalId=null, transactionTimeoutMs=2147483647, producerId=-1, producerEpoch=-1) with context RequestContext(header=RequestHeader(apiKey=INIT_PRODUCER_ID, apiVersion=4, clientId=producer-16, correlationId=2), connectionId='172.20.0.7:29092-172.20.0.1:43094-4', clientAddress=/172.20.0.1, principal=User:ANONYMOUS, listenerName=ListenerName(EXTERNAL), securityProtocol=PLAINTEXT, clientInformation=ClientInformation(softwareName=apache-kafka-java, softwareVersion=unknown), fromPrivilegedListener=false, principalSerde=Optional[org.apache.kafka.common.security.authenticator.DefaultKafkaPrincipalBuilder@4c7701d4]) (kafka.server.KafkaApis)
org.apache.kafka.common.errors.TimeoutException: Timed out waiting for next producer ID block
[2024-04-01 23:58:54,193] INFO [BrokerToControllerChannelManager broker=1 name=forwarding] Disconnecting from node 1 due to request timeout. (org.apache.kafka.clients.NetworkClient)
[2024-04-01 23:58:54,194] INFO [BrokerToControllerChannelManager broker=1 name=forwarding] Cancelled in-flight API_VERSIONS request with correlation id 1 due to node 1 being disconnected (elapsed time since creation: 30029ms, elapsed time since send: 30029ms, request timeout: 30000ms) (org.apache.kafka.clients.NetworkClient)
[2024-04-01 23:58:54,195] INFO [BrokerToControllerChannelManager broker=1 name=forwarding]: Recorded new controller, from now on will use broker nids-kafka-cntr:9092 (id: 1 rack: null) (kafka.server.BrokerToControllerRequestThread)

Zeek生产者日志

[disconnected broker msgs]
%4|1712015304.258|FAIL|rdkafka#producer-7| [thrd:localhost:29092/bootstrap]: localhost:29092/bootstrap: ApiVersionRequest failed: Local: Timed out: probably due to broker version < 0.10 (see api.version.request configuration) (after 10009ms in state APIVERSION_QUERY, 3 identical error(s) suppressed)
%3|1712015304.258|ERROR|rdkafka#producer-7| [thrd:localhost:29092/bootstrap]: 1/1 brokers are down
%4|1712015304.258|REQTMOUT|rdkafka#producer-7| [thrd:localhost:29092/bootstrap]: localhost:29092/bootstrap: Timed out 1 in-flight, 0 retry-queued, 0 out-queue, 0 partially-sent requests

Logstash消费者日志

[2024-04-01T23:58:31,613][INFO ][org.apache.kafka.clients.consumer.internals.AbstractCoordinator][main][b13590b3df4211af5265ef2746ad3101d589af4a46691ce4ad66a9091cfd09cb] [Consumer clientId=logstash-0, groupId=logstash] Group coordinator nids-kafka-cntr:9092 (id: 2147483646 rack: null) is unavailable or invalid, will attempt rediscovery
[2024-04-01T23:59:04,649][INFO ][org.apache.kafka.clients.FetchSessionHandler][main][b13590b3df4211af5265ef2746ad3101d589af4a46691ce4ad66a9091cfd09cb] [Consumer clientId=logstash-0, groupId=logstash] Error sending fetch request (sessionId=430766704, epoch=14) to node 1: {}.
org.apache.kafka.common.errors.DisconnectException: null

请问该问题的原因可能是什么?


内容的提问来源于stack exchange,提问作者Wyler Zahm

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.26 18:07:14