You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

使用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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.11 15:59:54