不同Docker容器内Flink作业向Kafka Topic写入数据失败排查求解
问题根因与解决方案
核心问题根因
- Flink作业运行在TaskManager容器内部,你配置的
bootstrap.servers = localhost:9092指向的是TaskManager容器本身的9092端口,而非Kafka容器的服务端口,自然无法建立连接 - Kafka的监听器配置存在绑定限制,原有配置中
PLAINTEXT_HOST绑定了localhost:9092,仅支持Kafka容器内部访问,其他同网络容器无法通过该地址通信 - TaskManager启动时多余映射了9092端口,会占用宿主机端口导致Kafka端口映射冲突
解决方案(兼容后续跨节点部署,非临时捷径)
步骤1:修改Flink作业Kafka配置
将作业中bootstrap地址改为Docker网络内Kafka的服务名+容器内通信端口,修改后代码如下:
val env: StreamExecutionEnvironment = StreamExecutionEnvironment.getExecutionEnvironment val transactions: DataStream[Transaction] = env .addSource(new TransactionSource) .name("transactions") val properties = new Properties // 修改为Docker网络内Kafka服务的访问地址 properties.setProperty("bootstrap.servers", "kafka:29092") val myProducer = new FlinkKafkaProducer[Transaction]( "flagged-transactions", // target topic TransactionSchema("flagged-transactions"), properties, // producer config FlinkKafkaProducer.Semantic.EXACTLY_ONCE) // fault-tolerance transactions.addSink(myProducer) env.execute("transactions")
同Docker网络内会自动解析kafka服务名对应Kafka容器的IP,29092是专门用于容器内部通信的监听器端口。
步骤2:修正docker-compose.yml Kafka配置
移除废弃参数,修正监听器绑定规则,修改后完整docker-compose.yml如下:
version: '2' networks: default: external: true name: flink-network services: zookeeper: image: confluentinc/cp-zookeeper:latest environment: ZOOKEEPER_CLIENT_PORT: 2181 ZOOKEEPER_TICK_TIME: 2000 ports: - 22181:2181 kafka: image: confluentinc/cp-kafka:latest depends_on: - zookeeper ports: # 宿主机本地访问Kafka用的端口映射 - 9092:9092 command: sh -c "((sleep 15 && kafka-topics --create --bootstrap-server localhost:29092 --replication-factor 1 --partitions 3 --topic flagged-transactions)&) && /etc/confluent/docker/run " environment: KAFKA_BROKER_ID: 1 KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181 # 监听器绑定0.0.0.0解除网卡访问限制 KAFKA_LISTENERS: PLAINTEXT://0.0.0.0:29092,PLAINTEXT_HOST://0.0.0.0:9092 KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://kafka:29092,PLAINTEXT_HOST://localhost:9092 KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT,PLAINTEXT_HOST:PLAINTEXT KAFKA_INTER_BROKER_LISTENER_NAME: PLAINTEXT KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
后续需要跨节点部署时,仅需要把KAFKA_ADVERTISED_LISTENERS中PLAINTEXT_HOST对应的localhost替换为宿主机的公网/内网IP即可,容器内部通信逻辑完全不需要调整。
步骤3:修正TaskManager启动命令
删除多余的9092端口映射,启动命令修改为:
docker run -d --rm --name=taskmanager --network flink-network --env FLINK_PROPERTIES="jobmanager.rpc.address: jobmanager" flink:latest taskmanager
步骤4:重新部署验证
- 停止所有已有容器:
docker stop jobmanager taskmanager,进入docker-compose目录执行docker-compose down - 重新启动Kafka:
docker-compose up -d,等待30秒后执行docker exec kafka kafka-topics --list --bootstrap-server localhost:29092确认flagged-transactions主题存在 - 重新启动Flink JobManager、TaskManager
- 重新打包Flink作业后提交到Flink Web UI即可正常运行
内容的提问来源于stack exchange,提问作者Zac-K
相关产品推荐
相关产品推荐

