为何Flink Kafka客户端配置连172.17.0.1:9092却尝试连接localhost:9092?
问题:Flink容器中Kafka客户端配置指定地址却尝试连接localhost:9092
我正尝试通过docker-compose搭建Flink JobManager-TaskManager集群,配置如下:
version: "3.7" services: jobmanagerconfig: image: flink:1.13.2-scala_2.12 expose: - "6133" - "6123" ports: - "8085:8081" command: standalone-job --job-classname net.mongerbot.configManager.App volumes: - ./usrlib/:/opt/flink/usrlib environment: - | FLINK_PROPERTIES= jobmanager.rpc.address: jobmanagerconfig parallelism.default: 2 taskmanager.numberOfTaskSlots: 4 - KAFKA_URI=${KAFKA_URI} - KAFKA_PORT=${KAFKA_PORT} - KAFKA_groupId=${KAFKA_groupId} taskmanagerconfig: image: flink:1.13.2-scala_2.12 depends_on: - jobmanagerconfig links: - jobmanagerconfig command: taskmanager # scale: 1 volumes: - ./usrlib/:/opt/flink/usrlib environment: - KAFKA_URI=${KAFKA_URI} - KAFKA_PORT=${KAFKA_PORT} - KAFKA_groupId=${KAFKA_groupId} - | FLINK_PROPERTIES= jobmanager.rpc.address: jobmanagerconfig parallelism.default: 2 taskmanager.numberOfTaskSlots: 4 volumes: usrlib: networks: default: external: name: mongerbot_network
两个容器中的环境变量值均正确,日志显示Kafka客户端已配置为连接172.17.0.1:9092:
docker-taskmanagerconfig-1 | 2022-12-08 09:36:56,065 INFO org.apache.kafka.clients.consumer.ConsumerConfig [] - ConsumerConfig values: docker-taskmanagerconfig-1 | allow.auto.create.topics = true docker-taskmanagerconfig-1 | auto.commit.interval.ms = 5000 docker-taskmanagerconfig-1 | auto.offset.reset = latest docker-taskmanagerconfig-1 | bootstrap.servers = [172.17.0.1:9092] docker-taskmanagerconfig-1 | check.crcs = true docker-taskmanagerconfig-1 | client.dns.lookup = default docker-taskmanagerconfig-1 | client.id = docker-taskmanagerconfig-1 | client.rack = docker-taskmanagerconfig-1 | connections.max.idle.ms = 540000 docker-taskmanagerconfig-1 | default.api.timeout.ms = 60000 docker-taskmanagerconfig-1 | enable.auto.commit = true docker-taskmanagerconfig-1 | exclude.internal.topics = true ...
但紧随其后的日志显示客户端尝试连接localhost:9092:
docker-taskmanagerconfig-1 | value.deserializer = class org.apache.kafka.common.serialization.ByteArrayDeserializer docker-taskmanagerconfig-1 | docker-taskmanagerconfig-1 | 2022-12-08 09:36:56,084 INFO org.apache.kafka.clients.consumer.KafkaConsumer [] - [Consumer clientId=consumer-configManager-7, groupId=configManager] Subscribed to partition(s): config.subscribe-0, config.subscribe-2 docker-taskmanagerconfig-1 | 2022-12-08 09:36:56,090 INFO org.apache.kafka.clients.Metadata [] - [Consumer clientId=consumer-configManager-7, groupId=configManager] Cluster ID: s2iVODWcQ2Kbw4R5jL6RCw docker-taskmanagerconfig-1 | 2022-12-08 09:36:56,091 INFO org.apache.kafka.clients.consumer.internals.AbstractCoordinator [] - [Consumer clientId=consumer-configManager-7, groupId=configManager] Discovered group coordinator localhost:9092 (id: 2147483646 rack: null) docker-taskmanagerconfig-1 | 2022-12-08 09:36:56,094 WARN org.apache.kafka.clients.NetworkClient [] - [Consumer clientId=consumer-configManager-7, groupId=configManager] Connection to node 2147483646 (localhost/127.0.0.1:9092) could not be established. Broker may not be available. docker-taskmanagerconfig-1 | 2022-12-08 09:36:56,094 INFO org.apache.kafka.clients.consumer.internals.AbstractCoordinator [] - [Consumer clientId=consumer-configManager-7, groupId=configManager] Group coordinator localhost:9092 (id: 2147483646 rack: null) is unavailable or invalid, will attempt rediscovery docker-taskmanagerconfig-1 | 2022-12-08 09:36:56,094 INFO org.apache.kafka.common.utils.AppInfoParser [] - Kafka version: 2.4.1 docker-taskmanagerconfig-1 | 2022-12-08 09:36:56,095 INFO org.apache.kafka.common.utils.AppInfoParser [] - Kafka commitId: c57222ae8cd7866b docker-taskmanagerconfig-1 | 2022-12-08 09:36:56,095 INFO org.apache.kafka.common.utils.AppInfoParser [] - Kafka startTimeMs: 1670492216094 docker-taskmanagerconfig-1 | 2022-12-08 09:36:56,096 INFO org.apache.kafka.clients.consumer.KafkaConsumer [] - [Consumer clientId=consumer-configManager-8, groupId=configManager] Subscribed to partition(s): config.subscribe-1 docker-taskmanagerconfig-1 | 2022-12-08 09:36:56,101 INFO org.apache.kafka.clients.Metadata [] - [Consumer clientId=consumer-configManager-8, groupId=configManager] Cluster ID: s2iVODWcQ2Kbw4R5jL6RCw docker-taskmanagerconfig-1 | 2022-12-08 09:36:56,102 INFO org.apache.kafka.clients.consumer.internals.AbstractCoordinator [] - [Consumer clientId=consumer-configManager-8, groupId=configManager] Discovered group coordinator localhost:9092 (id: 2147483646 rack: null) docker-taskmanagerconfig-1 | 2022-12-08 09:36:56,103 WARN org.apache.kafka.clients.NetworkClient [] - [Consumer clientId=consumer-configManager-8, groupId=configManager] Connection to node 2147483646 (localhost/127.0.0.1:9092) could not be established. Broker may not be available. docker-taskmanagerconfig-1 | 2022-12-08 09:36:56,104 INFO org.apache.kafka.clients.consumer.internals.AbstractCoordinator [] - [Consumer clientId=consumer-configManager-8, groupId=configManager] Group coordinator localhost:9092 (id: 2147483646 rack: null) is unavailable or invalid, will attempt rediscovery docker-taskmanagerconfig-1 | 2022-12-08 09:36:56,197 WARN org.apache.kafka.clients.NetworkClient [] - [Consumer clientId=consumer-configManager-7, groupId=configManager] Connection to node 1 (localhost/127.0.0.1:9092) could not be established. Broker may not be available. docker-taskmanagerconfig-1 | 2022-12-08 09:36:56,207 WARN org.apache.kafka.clients.NetworkClient
原因分析
这是Kafka的典型配置问题:Kafka Broker的advertised.listeners参数配置错误。
Kafka的工作逻辑是:客户端先通过bootstrap.servers配置的地址连接到Broker,Broker会将自己配置的advertised.listeners地址返回给客户端,后续客户端所有的通信(包括连接Group Coordinator、发送/接收消息)都会使用这个返回的地址,而非初始的bootstrap地址。
如果你的Kafka Broker的advertised.listeners设置为localhost:9092,那么即使客户端一开始用172.17.0.1:9092连接,后续也会被引导去连接容器内的localhost(也就是容器自身),而容器内并没有运行Kafka,因此会出现连接失败的报错。
解决办法
- 找到Kafka的配置文件(通常名为
server.properties) - 修改
advertised.listeners配置:- 如果Kafka运行在宿主机上,设置为宿主机的可访问IP(比如
172.17.0.1:9092,和你客户端bootstrap.servers用的地址一致) - 如果Kafka也运行在Docker容器中,且和Flink集群在同一个Docker网络下,设置为Kafka容器的服务名(比如
kafka:9092)
- 如果Kafka运行在宿主机上,设置为宿主机的可访问IP(比如
- 重启Kafka服务,让配置生效
内容的提问来源于stack exchange,提问作者Mohammad Esmaeilzadeh
相关产品推荐
相关产品推荐

