Docker Compose环境中Spark连接Kafka失败问题求助
问题:Spark无法连接Docker Compose部署的Kafka服务
执行Spark消费Kafka数据的作业时,出现以下错误:
24/12/13 23:28:02 WARN NetworkClient: [AdminClient clientId=adminclient-1] Connection to node 1 (localhost/127.0.0.1:9092) could not be established. Broker may not be available.
部署配置与命令
docker-compose.yml配置
version: '3.8' services: zookeeper: image: confluentinc/cp-zookeeper:latest environment: ZOOKEEPER_CLIENT_PORT: 2181 ZOOKEEPER_TICK_TIME: 2000 kafka: image: confluentinc/cp-kafka:latest ports: - "9092:9092" environment: KAFKA_BROKER_ID: 1 KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181 KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://localhost:9092 KAFKA_LISTENERS: PLAINTEXT://0.0.0.0:9092 spark: image: bitnami/spark:latest ports: - "8080:8080" volumes: - 'C:/docker-proyecto/spark-data:/opt/spark' environment: - SPARK_MODE=master - SPARK_RPC_AUTHENTICATION_ENABLED=no spark-worker: image: bitnami/spark:latest environment: - SPARK_MODE=worker - SPARK_MASTER_URL=spark://spark-master:7077 - SPARK_WORKER_MEMORY=1G - SPARK_WORKER_CORES=1 depends_on: - spark volumes: - 'C:/docker-proyecto/spark-data:/opt/spark'
执行的spark-submit命令
spark-submit --packages org.apache.spark:spark-sql-kafka-0-10_2.12:3.5.3 /opt/spark/script.py
Spark脚本内容
kafka_servers = "localhost:9092" stream_df = spark.readStream \ .format("kafka") \ .option("kafka.bootstrap.servers", kafka_servers) \ .option("subscribe", "topic1,topic2") \ .option("startingOffsets", "earliest") \ .load()
预期目标:Spark正常连接Kafka读取Topic数据,处理后存入HDFS。
问题原因与解决方法
1. 容器间网络访问错误
Spark容器内的localhost指向容器自身,而非宿主机,所以用localhost:9092无法访问到Kafka容器。Docker Compose会自动为服务创建内部DNS,直接使用Kafka服务名kafka即可访问对应容器。
2. Kafka监听地址配置缺陷
当前KAFKA_ADVERTISED_LISTENERS仅配置了宿主机可访问的地址,当Spark容器从Kafka获取broker地址时,得到的是localhost:9092,导致Spark容器尝试连接自己的9092端口,必然失败。需要为Kafka配置两个监听地址:一个供宿主机外部访问,一个供Docker内部容器间访问。
修改步骤
步骤1:更新docker-compose.yml的Kafka配置
修改Kafka服务的环境变量,添加内部监听地址:
kafka: image: confluentinc/cp-kafka:latest ports: - "9092:9092" - "9093:9093" # 新增端口映射供内部访问 environment: KAFKA_BROKER_ID: 1 KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181 # 配置两个advertised listeners:外部用localhost,内部用服务名 KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://localhost:9092,PLAINTEXT_INTERNAL://kafka:9093 KAFKA_LISTENERS: PLAINTEXT://0.0.0.0:9092,PLAINTEXT_INTERNAL://0.0.0.0:9093 KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT,PLAINTEXT_INTERNAL:PLAINTEXT KAFKA_INTER_BROKER_LISTENER_NAME: PLAINTEXT_INTERNAL
步骤2:修改Spark脚本的Kafka地址
将脚本中的kafka_servers改为Docker内部可访问的地址:
kafka_servers = "kafka:9093" # 使用Kafka服务名+内部监听端口 stream_df = spark.readStream \ .format("kafka") \ .option("kafka.bootstrap.servers", kafka_servers) \ .option("subscribe", "topic1,topic2") \ .option("startingOffsets", "earliest") \ .load()
步骤3:验证连接
- 重启Docker Compose服务:
docker-compose down && docker-compose up -d - 进入Spark容器测试网络连通性:
docker exec -it <spark-container-id> ping kafka,能ping通则内部网络正常 - 可在Kafka容器内创建测试Topic:
docker exec -it <kafka-container-id> kafka-topics.sh --create --topic topic1 --bootstrap-server localhost:9092 --partitions 1 --replication-factor 1
这样修改后,Spark容器就能通过内部网络正常连接到Kafka服务,同时宿主机依然可以通过localhost:9092访问Kafka。
内容的提问来源于stack exchange,提问作者Leynder Sánchez
相关产品推荐
相关产品推荐

