升级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
相关产品推荐
相关产品推荐

