Docker环境下Spark Structured Streaming无法读取Kafka问题求助
故障排查与解决方案
根因分析
报错核心为java.net.UnknownHostException: kafka,本质是运行Kafka消费者的进程无法解析kafka主机名:
- Spark Structured Streaming的Kafka消费者运行在Executor节点,而非Driver节点
- 此前Spark Streaming正常通常是因为作业以local模式运行在提交任务的容器中,而当前作业提交到Spark集群后,实际运行在Worker节点,需要Worker节点能正常解析Kafka主机名并连通
排查步骤
1. 验证Spark Worker网络连通性
进入任意一个Spark Worker容器执行命令,验证网络可达:
# 验证域名解析 ping kafka # 验证端口连通 telnet kafka 9092
如果ping不通,说明Spark Worker未正确接入twitter-streaming_default网络,检查第二份docker-compose配置中是否给spark-worker服务单独指定了其他网络,若单独指定则不会默认加入公共外部网络。
2. 修正Kafka配置
你使用的KAFKA_ADVERTISED_HOST_NAME是Kafka 0.9.x版本的废弃参数,新版本Kafka已不识别该配置,需替换为标准监听配置,修改第一份docker-compose的kafka服务环境变量:
kafka: image: linuxkitpoc/kafka:latest ports: - "9092:9092" depends_on: - zookeeper environment: KAFKA_LISTENERS: PLAINTEXT://0.0.0.0:9092 KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://kafka:9092 KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181 KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT KAFKA_INTER_BROKER_LISTENER_NAME: PLAINTEXT volumes: - /var/run/docker.sock:/var/run/docker.sock
修改后重启Kafka容器生效。
3. 修正Structured Streaming代码冗余参数
Spark Structured Streaming的Kafka Source会自动管理消费位点,不支持用户自定义group.id和enable.auto.commit参数,这两个配置会被Spark自动覆盖,属于无效配置,可直接删除避免兼容性问题:
val df = spark .readStream .format("kafka") .option("kafka.bootstrap.servers", "kafka:9092") .option("startingOffsets", "earliest") .option("subscribe", "tweet-upload-6") .option("failOnDataLoss", false) .load()
4. 临时兜底方案(域名解析异常时使用)
如果确认网络已连通但域名解析始终失败,可通过固定IP绑定的方式解决:
- 先查询Kafka容器的实际IP:
docker inspect -f '{{range.NetworkSettings.Networks}}{{.IPAddress}}{{end}}' kafka容器名
- 两种绑定方式二选一:
- 直接将代码中
kafka:9092替换为查询到的IP,比如172.18.0.4:9092 - 在第二份docker-compose的spark-worker服务配置中添加host绑定:
spark-worker-1: # 原有配置保留 extra_hosts: - "kafka:172.18.0.4" spark-worker-2: # 原有配置保留 extra_hosts: - "kafka:172.18.0.4"
将172.18.0.4替换为你实际查询到的Kafka容器IP即可。
内容的提问来源于stack exchange,提问作者Milisav Kovacevic
相关产品推荐
相关产品推荐

