使用Spark Structured Streaming预览Kafka数据时PySpark启动卡住求助
问题描述
尝试通过Spark Structured Streaming从Kafka主题消费数据,运行PySpark命令后进程卡住数分钟,无法进入Spark CLI。已确认Kafka主题有数据写入,执行命令如下:
pyspark \ --master yarn \ --packages org.apache.spark:spark-sql-kafka-0-10_2.12:3.3.3 \ --conf spark.ui.port=0 \ --conf spark.sql.warehouse.dir=/user/${USER}/warehouse
输出显示依赖包已成功下载,但最后停留在重复资源添加到分布式缓存的警告后,无后续进展。
解决方案
1. 处理重复资源缓存警告
重复添加依赖到分布式缓存可能导致YARN容器启动卡住,可通过两种方式解决:
- 手动指定本地依赖包:直接用
--jars参数指定已下载的本地jar包路径,避免Ivy重复解析:pyspark \ --master yarn \ --jars /home/dee/.ivy2/jars/org.apache.spark_spark-sql-kafka-0-10_2.12-3.3.3.jar,/home/dee/.ivy2/jars/org.apache.spark_spark-token-provider-kafka-0-10_2.12-3.3.3.jar,/home/dee/.ivy2/jars/org.apache.kafka_kafka-clients-2.8.1.jar,/home/dee/.ivy2/jars/com.google.code.findbugs_jsr305-3.0.0.jar,/home/dee/.ivy2/jars/org.apache.commons_commons-pool2-2.11.1.jar,/home/dee/.ivy2/jars/org.spark-project.spark_unused-1.0.0.jar,/home/dee/.ivy2/jars/org.apache.hadoop_hadoop-client-runtime-3.3.2.jar,/home/dee/.ivy2/jars/org.lz4_lz4-java-1.8.0.jar,/home/dee/.ivy2/jars/org.xerial.snappy_snappy-java-1.1.8.4.jar,/home/dee/.ivy2/jars/org.slf4j_slf4j-api-1.7.32.jar,/home/dee/.ivy2/jars/org.apache.hadoop_hadoop-client-api-3.3.2.jar,/home/dee/.ivy2/jars/commons-logging_commons-logging-1.1.3.jar \ --conf spark.ui.port=0 \ --conf spark.sql.warehouse.dir=/user/${USER}/warehouse - 禁用重复资源检查:在命令中添加配置跳过重复资源警告:
pyspark \ --master yarn \ --packages org.apache.spark:spark-sql-kafka-0-10_2.12:3.3.3 \ --conf spark.ui.port=0 \ --conf spark.sql.warehouse.dir=/user/${USER}/warehouse \ --conf spark.yarn.dist.duplicateWarning=false
2. 检查YARN集群资源
进程卡住可能是YARN无足够资源分配给Spark组件:
- 执行
yarn node -list查看集群可用节点状态 - 执行
yarn application -list查看当前运行的应用,确认资源占用情况 - 若资源不足,可降低Spark资源配置后重试:
pyspark \ --master yarn \ --packages org.apache.spark:spark-sql-kafka-0-10_2.12:3.3.3 \ --conf spark.ui.port=0 \ --conf spark.sql.warehouse.dir=/user/${USER}/warehouse \ --conf spark.driver.memory=1g \ --conf spark.executor.memory=1g \ --conf spark.executor.cores=1
3. 查看详细日志定位问题
通过YARN日志工具获取卡住的具体原因:
- 用
yarn application -list找到当前Spark应用的ID - 执行
yarn logs -applicationId <你的应用ID>查看完整日志
4. 本地模式验证环境
先切换到本地模式运行,排除环境依赖问题:
pyspark \ --master local[*] \ --packages org.apache.spark:spark-sql-kafka-0-10_2.12:3.3.3
若本地模式能正常进入CLI,说明问题出在YARN集群配置或资源上。
内容的提问来源于stack exchange,提问作者Oyindamola Victor
相关产品推荐
相关产品推荐

