Docker环境下Spark集群连接Kafka的配置方案咨询
如何让独立Docker环境中的Spark集群读取Kafka主题数据?
针对你的场景,需要完成以下几个关键配置步骤,实现Spark与Kafka的跨Docker环境连通:
1. 统一Docker网络,实现跨容器互通
两个独立的docker-compose默认使用各自的隔离网络,需将Kafka和Spark集群加入同一个自定义网络:
- 创建自定义Docker网络:
docker network create streaming-poc-network
- 修改Kafka的
docker-compose.yml,在最外层添加网络配置(原有服务配置保留):
version: '2' networks: default: external: name: streaming-poc-network services: # 原有zookeeper、kafka、mysql、connect配置不变
- 修改Spark的
docker-compose.yml,同样添加网络配置(原有服务配置保留):
version: '3' networks: default: external: name: streaming-poc-network services: # 原有spark-master、spark-worker配置不变
- 重启两个环境的容器,确保它们成功加入同一网络。
2. 调整Kafka监听配置,允许外部容器访问
Debezium Kafka默认仅支持容器内部访问,需修改Kafka的环境变量,添加对外监听规则:
在Kafka的docker-compose.yml的kafka服务中,新增以下环境变量:
environment: # 原有ZOOKEEPER_CONNECT配置保留 - LISTENERS=PLAINTEXT://0.0.0.0:9092,PLAINTEXT_HOST://0.0.0.0:9093 - ADVERTISED_LISTENERS=PLAINTEXT://kafka:9092,PLAINTEXT_HOST://localhost:9092 - LISTENER_SECURITY_PROTOCOL_MAP=PLAINTEXT:PLAINTEXT,PLAINTEXT_HOST:PLAINTEXT
说明:
kafka:9092供同一Docker网络内的服务(如Debezium Connect、Spark容器)访问localhost:9092供宿主机或跨网络容器通过宿主机端口访问
3. 配置Spark任务的Kafka连接参数
提交Spark流处理任务时,需指定正确的Kafka地址并引入依赖:
3.1 设置正确的bootstrap服务器地址
在Spark代码中,根据部署场景配置bootstrap.servers:
- 如果Spark容器在同一自定义网络内,使用
kafka:9092 - 如果从宿主机提交任务到Spark集群,使用
localhost:9092
示例Scala代码片段:
val df = spark.readStream .format("kafka") .option("kafka.bootstrap.servers", "kafka:9092") // 同一网络内用此地址,否则替换为localhost:9092 .option("subscribe", "your-debezium-topic") // 替换为你的Debezium同步主题 .load()
3.2 引入Kafka连接器依赖
Spark 3.3需要对应版本的Kafka连接器jar包,提交任务时通过--packages参数引入:
spark-submit \ --master spark://spark-master:7077 \ --packages org.apache.spark:spark-sql-kafka-0-10_2.12:3.3.0 \ --class com.your.package.YourStreamingApp \ your-app.jar
注意:依赖版本需与Spark版本严格匹配,Scala版本(此处为2.12)也要对应Spark默认版本。
4. 验证连通性
- 进入Spark主节点容器,执行
telnet kafka 9092测试端口连通性 - 提交测试任务,检查Spark是否能成功读取Kafka主题的数据
内容的提问来源于stack exchange,提问作者Eugene Goldberg
相关产品推荐
相关产品推荐

