Spark Streaming写入Kafka超时:Topic RANDOM_STRING元数据获取失败
核心原因定位
从报错TimeoutException: Topic RANDOM_STRING not present in metadata after 60000 ms和日志中executor 1: server.com来看,最可能的问题是Spark Executor无法正确连接到Kafka Broker,导致无法获取主题元数据。控制台输出正常说明流读取逻辑无问题,问题集中在Kafka的网络配置或Spark的连接配置上。
具体解决方案
1. 修正Kafka的监听地址配置
当前server.properties中listeners=PLAINTEXT://:9092仅绑定了端口,但未指定对外暴露的地址。当Spark Executor运行在另一台机器(如日志中的server.com)时,Kafka返回的元数据中Broker地址是localhost:9092,而Executor的localhost指向自身而非Kafka所在机器,导致连接超时。
修改server.properties:
# 保持监听所有网卡端口 listeners=PLAINTEXT://0.0.0.0:9092 # 指定Executor能访问的Kafka实际地址(替换为Kafka所在机器的公网/内网IP) advertised.listeners=PLAINTEXT://192.168.x.x:9092
修改后重启Kafka服务。
2. 更新Spark代码中的Kafka连接地址
将kafka.bootstrap.servers从localhost:9092改为Kafka所在机器的实际IP,确保Executor能正确解析到Broker:
kafka_config = { "kafka.bootstrap.servers": "192.168.x.x:9092", # 替换为实际IP "checkpointLocation": "/user/aiman/checkpoint/kafka_local/random_string", "topic": "RANDOM_STRING" }
3. 验证网络连通性
在Executor所在的server.com机器上,执行以下命令测试Kafka端口是否可达:
nc -zv 192.168.x.x 9092
若连接失败,检查Kafka所在机器的防火墙规则,开放9092端口。
4. 确保DataFrame结构符合Kafka连接器要求
Spark SQL Kafka连接器要求写入的DataFrame必须包含value列(字符串或二进制类型)。Socket源读取的DataFrame默认包含value列(字符串类型),可显式确认类型避免潜在问题:
from pyspark.sql.functions import col data = data.select(col("value").cast("string"))
5. 检查Kafka主题状态
确认主题的副本和分区状态正常:
./kafka-topics.sh --zookeeper localhost:2181 --describe --topic RANDOM_STRING
确保Isr列包含Leader节点,说明主题处于可用状态。
版本兼容性说明
你使用的spark-sql-kafka-0-10_2.11-2.4.0-cdh6.3.4.jar兼容Kafka 0.10.0.0至2.0.0版本,因此更换Kafka版本至0.10.0.0并非问题根源,无需调整版本。
内容的提问来源于stack exchange,提问作者aiman

