Spark独立模式下DataFrame保存为Parquet失败及临时目录问题
我在搭载Hive Catalog的Spark独立模式下运行任务,尝试从外部文档加载数据并将其以Parquet格式保存到磁盘。执行的代码如下:
rdd = sc \ .textFile('/data/source.txt', NUM_SLICES) \ .map(lambda x: (x[:5], x[6:12], gensim.utils.simple_preprocess(x[13:]))) schema = StructType([ StructField('c1', StringType(), False), StructField('c2', StringType(), False), StructField('c3', ArrayType(StringType(), True), False), ]) data = sql_context.createDataFrame(rdd, schema) data.write.mode('overwrite').parquet('/data/some_dir')
但读取该文件时失败,抛出异常:
AnalysisException: 'Unable to infer schema for Parquet. It must be specified manually.;'
查看所有worker节点的存储路径后,发现文件结构如下(通过clush -ab 'locate some_file'获取):
--------------- master --------------- /data/some_file /data/some_file/._SUCCESS.crc /data/some_file/_SUCCESS --------------- worker1 --------------- /data/some_file /data/some_file/_temporary /data/some_file/_temporary/0 /data/some_file/_temporary/0/_temporary /data/some_file/_temporary/0/task_20180511211832_0010_m_000000 /data/some_file/_temporary/0/task_20180511211832_0010_m_000039 /data/some_file/_temporary/0/task_20180511211832_0010_m_000000/.part-00000-1b2764a6-28a3-4ba2-9493-766074eef4d5-c000.snappy.parquet.crc /data/some_file/_temporary/0/task_20180511211832_0010_m_000000/part-00000-1b2764a6-28a3-4ba2-9493-766074eef4d5-c000.snappy.parquet /data/some_file/_temporary/0/task_20180511211832_0010_m_000039/.part-00039-1b2764a6-28a3-4ba2-9493-766074eef4d5-c000.snappy.parquet.crc /data/some_file/_temporary/0/task_20180511211832_0010_m_000039/part-00039-1b2764a6-28a3-4ba2-9493-766074eef4d5-c000.snappy.parquet --------------- worker2 --------------- /data/some_file /data/some_file/_temporary /data/some_file/_temporary/0 /data/some_file/_temporary/0/_temporary /data/some_file/_temporary/0/task_20180511211832_0010_m_000011 /data/some_file/_temporary/0/task_20180511211832_0010_m_000017 /data/some_file/_temporary/0/task_20180511211832_0010_m_000029 /data/some_file/_temporary/0/task_20180511211832_0010_m_000038 /data/some_file/_temporary/0/task_20180511211832_0010_m_000011/.part-00011-1b2764a6-28a3-4ba2-9493-766074eef4d5-c000.snappy.parquet.crc /data/some_file/_temporary/0/task_20180511211832_0010_m_000011/part-00011-1b2764a6-28a3-4ba2-9493-766074eef4d5-c000.snappy.parquet /data/some_file/_temporary/0/task_20180511211832_0010_m_000017/.part-00017-1b2764a6-28a3-4ba2-9493-766074eef4d5-c000.snappy.parquet.crc /data/some_file/_temporary/0/task_20180511211832_0010_m_000017/part-00017-1b2764a6-28a3-4ba2-9493-766074eef4d5-c000.snappy.parquet /data/some_file/_temporary/0/task_20180511211832_0010_m_000029/.part-00029-1b2764a6-28a3-4ba2-9493-766074eef4d5-c000.snappy.parquet.crc /data/some_file/_temporary/0/task_20180511211832_0010_m_000029/part-00029-1b2764a6-28a3-4ba2-9493-766074eef4d5-c000.snappy.parquet /data/some_file/_temporary/0/task_20180511211832_0010_m_000038/.part-00038-1b2764a6-28a3-4ba2-9493-766074eef4d5-c000.snappy.parquet.crc /data/some_file/_temporary/0/task_20180511211832_0010_m_000038/part-00038-1b2764a6-28a3-4ba2-9493-766074eef4d5-c000.snappy.parquet --------------- worker3 --------------- /data/some_file /data/some_file/_temporary /data/some_file/_temporary/0 /data/some_file/_temporary/0/_temporary /data/some_file/_temporary/0/task_20180511211832_0010_m_000040 /data/some_file/_temporary/0/task_20180511211832_0010_m_000043 /data/some_file/_temporary/0/task_20180511211832_0010_m_000046 /data/some_file/_temporary/0/task_20180511211832_0010_m_000040/.part-00040-1b2764a6-28a3-4ba2-9493-766074eef4d5-c000.snappy.parquet.crc /data/some_file/_temporary/0/task_20180511211832_0010_m_000040/part-00040-1b2764a6-28a3-4ba2-9493-766074eef4d5-c000.snappy.parquet /data/some_file/_temporary/0/task_20180511211832_0010_m_000043/.part-00043-1b2764a6-28a3-4ba2-9493-766074eef4d5-c000.snappy.parquet.crc /data/some_file/_temporary/0/task_20180511211832_0010_m_000043/part-00043-1b2764a6-28a3-4ba2-9493-766074eef4d5-c000.snappy.parquet /data/some_file/_temporary/0/task_20180511211832_0010_m_000046/.part-00046-1b2764a6-28a3-4ba2-9493-766074eef4d5-c000.snappy.parquet.crc /data/some_file/_temporary/0/task_20180511211832_0010_m_000046/part-00046-1b2764a6-28a3-4ba2-9493-766074eef4d5-c000.snappy.parquet
问题原因分析
Spark在写入Parquet(或其他分布式存储格式)时,会遵循以下流程:
- 每个Worker任务先将数据写入本地临时目录(
_temporary下的子目录) - 当所有Worker任务完成后,Driver节点会将所有临时目录中的数据文件移动到最终目标目录
- 最后Driver删除临时目录,并生成
_SUCCESS文件标记任务完成
现在你的场景中,临时目录没有被清理,最终目录下只有_SUCCESS相关文件,没有实际的Parquet数据文件,导致读取时无法推断Schema,抛出异常。核心原因是:
- 使用了本地文件系统而非分布式存储:Spark独立模式下,如果每个节点的
/data是本地磁盘(而非HDFS、S3等共享存储),Driver在Master节点无法访问Worker节点本地的临时文件,因此无法完成文件移动和临时目录清理的步骤。 - 任务可能异常终止:如果Driver或某个Worker任务意外退出,也会导致后续的提交步骤无法执行,临时目录残留。
解决方案
切换到分布式文件系统(推荐)
把存储路径改为分布式文件系统的路径(比如HDFS的hdfs://master:9000/data/some_dir),这样所有节点共享存储,Driver可以统一处理文件的移动和清理,这是Spark分布式任务的标准配置。检查任务执行日志
查看Driver日志和Worker节点的日志,确认是否有任务失败、Driver崩溃或其他异常信息,排查任务终止的原因。手动修复(临时测试方案)
如果是测试环境,可以手动将各个Worker节点_temporary目录下的Parquet文件移动到/data/some_file目录,然后删除所有节点的_temporary目录,之后再尝试读取:# 在每个Worker节点执行 mv /data/some_file/_temporary/0/task_*/part-*.snappy.parquet /data/some_file/ rm -rf /data/some_file/_temporary/
内容的提问来源于stack exchange,提问作者kirylm

