Spark-Submit运行Kafka流WordCount时连接失败及日志问题排查
从你提供的日志和执行命令来看,核心问题是Kafka Broker地址的格式不符合Spark Streaming的要求,导致无法正常建立连接。
问题根源
你执行的启动命令是:
spark-submit --packages org.apache.spark:spark-streaming-kafka-0-8_2.11:2.2.0 direct_kafka_wordcount.py localhost 9092
这里你把localhost(主机)和9092(端口)拆成了两个独立参数传给脚本,但direct_kafka_wordcount.py脚本期望的Broker地址格式是host:port的完整字符串,而非拆分的主机和端口。日志里的关键报错直接点明了这一点:
org.apache.spark.SparkException: Broker格式不正确,应为:[localhost]
这说明脚本接收到的Broker参数是单独的localhost(缺少端口)和9092(缺少主机),两者都不符合host:port的标准格式,因此触发了格式校验错误。
快速修复方案
方案1:修改命令行参数(最简单直接)
把Broker地址合并成localhost:9092作为单个参数传入,修改后的命令如下:
spark-submit --packages org.apache.spark:spark-streaming-kafka-0-8_2.11:2.2.0 direct_kafka_wordcount.py localhost:9092
这样脚本就能接收到正确格式的Broker地址,正常连接Kafka集群。
方案2:修改脚本代码(保留原命令参数格式)
如果你希望继续使用localhost 9092的参数拆分方式,可以修改direct_kafka_wordcount.py的参数处理逻辑:
找到脚本开头获取参数的代码,将传入的主机和端口拼接成标准格式:
# 原参数处理逻辑(示例) brokers = sys.argv[1] topic = sys.argv[2] # 修改为: broker_host = sys.argv[1] broker_port = sys.argv[2] brokers = f"{broker_host}:{broker_port}" topic = sys.argv[3] # 注意topic的参数索引需要对应后移一位
修改后再执行原命令即可正常运行。
额外兼容性提示
从日志可以看到你的Spark版本是2.4.3,但指定的Kafka Streaming包版本是2.2.0,版本差异可能带来潜在的兼容性问题。建议使用与Spark版本匹配的依赖包,比如:
--packages org.apache.spark:spark-streaming-kafka-0-8_2.11:2.4.3
这能避免因版本不兼容导致的其他未知问题。
内容的提问来源于stack exchange,提问作者sagar pawar

