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

Dataproc Pyspark读取GCS做地理空间转换后写入GCS分区卡住求助

问题排查与解决方案

1. 首要问题:分区列名笔误

你在执行写入分区的代码中,指定的分区列prov不存在:

# 错误写法
final_df.write.partitionBy('prov','year','month','day')

你代码中实际生成的省份分区列是province,此处属于笔误,修正为:

# 正确写法
final_df.write.partitionBy('province','year','month','day')

该错误会导致Spark在执行分区写入时无法找到对应列,任务逻辑阻塞。

2. 依赖版本冲突问题

你在两处重复指定依赖,且版本不匹配、存在无效依赖,会引发类加载异常导致任务卡住:

  • Airflow的Dataproc作业提交参数中已经指定了完整的jar包依赖,Spark代码中不需要再通过spark.jars.packages重复声明,且两处的Sedona、geotools版本不一致,会产生冲突
  • Airflow的jar列表中同时引入了hadoop3和hadoop2两个版本的gcs连接器,Dataproc 2.0版本使用Hadoop 3,需要删除gs://bucket/data/gcs-connector-hadoop2-2.1.3-shaded.jar这个无效依赖
    解决方法:删除Spark代码中SparkSession构造时所有spark.jars.packages相关的配置,同时清理Airflow中重复的无效依赖包。

3. Spark运行配置无效覆盖

SparkSession初始化完成后,再通过conf.set修改executor内存、driver内存等运行参数是不生效的,你代码中还存在配置名错误、单位缺失问题:

# 以下配置全部无效,且存在错误
conf.set("spark.executor.memory", "35g")
conf.set("spark.driver.memory", "35") # 缺失单位G
conf.set("spark.spark.executor.cores", '3') # 配置名多了一个spark前缀

所有Spark运行参数统一在Airflow创建集群的spark_config中配置即可,删除Spark代码中这几行无效配置。

4. 空间Join未优化

你当前使用的ST_Contains关联默认走笛卡尔积计算,哪怕小数据量也可能出现运行缓慢:

  • 省份表数据量很小,可以在SQL中添加广播提示强制广播小表,避免Shuffle:
SELECT /*+ BROADCAST(b) */ a.*, b.name as province
FROM raw_data_table as a
LEFT OUTER JOIN bq_provs as b
ON ST_Contains(b.geom, a.pt_geom)
  • 后续处理全量数据时,建议开启Sedona空间索引优化关联性能。

5. 日志查看方式错误

你当前看到的是Hadoop JobHistory服务的常规INFO日志,不属于Spark作业的运行日志,无法定位错误。你需要在Dataproc集群页面的作业标签页,找到对应Spark作业的Driver日志、Executor日志,查看具体的错误堆栈、任务运行状态。

后续全量数据处理建议

  • 按天分批提交作业,不要一次性加载4年全量数据,降低集群压力
  • 写入Parquet前添加repartition或coalesce控制输出文件大小,避免产生大量小文件
  • 开启Dynamic Partition Overwrite功能,避免重写整个表的所有分区。

内容的提问来源于stack exchange,提问作者cpatino08

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.27 18:24:02