如何使用PySpark读取变量列表内的Parquet文件并合并为单文件
实现方案
原有代码存在的4个前置问题
- 缺失
os模块导入,现有代码调用os.walk等接口会直接报名称错误 os系列是本地文件系统API,无法识别dbfs:/开头的协议路径,直接遍历会找不到目标目录,需替换为Databricks本地挂载的/dbfs/前缀路径filter()返回的是迭代器对象,不是实体列表,后续传入读取接口时容易出现遍历耗尽、类型不匹配问题,需要显式转为list类型- 未做文件后缀校验,遍历过程中会把目录下的
_SUCCESS标记文件、隐藏临时文件等非Parquet文件纳入筛选,导致后续读取出错
完整实现代码
import os import datetime from datetime import datetime from pyspark.sql import SparkSession # 本地调试时手动初始化SparkSession,Databricks环境可省略该段,直接使用内置的spark对象 spark = SparkSession.builder.appName("merge_filtered_parquet").getOrCreate() # os接口操作文件属性必须用/dbfs前缀的本地挂载路径 source_dir = '/dbfs/mnt/abc/def/efg' # 替换为你的目标输出路径,Spark写入支持dbfs:/协议 target_path = 'dbfs:/mnt/abc/def/merged_single_file.parquet' current_utc_time = datetime.utcnow() all_files = [] filtered_parquet_files = [] # 递归遍历目录收集所有文件路径 for dirpath, _, filenames in os.walk(source_dir): all_files.extend([os.path.join(dirpath, f) for f in filenames]) # 筛选修改时间早于当前UTC时间、后缀为.parquet的文件 filtered_parquet_files = list( filter( lambda x: x.endswith('.parquet') and datetime.utcfromtimestamp(os.path.getmtime(x)) < current_utc_time, all_files ) ) if len(filtered_parquet_files) == 0: raise RuntimeError("未筛选到符合条件的Parquet文件,终止执行") # 路径转换:将本地/dbfs路径换回Spark兼容的dbfs:/协议路径 spark_readable_paths = [p.replace('/dbfs/', 'dbfs:/') for p in filtered_parquet_files] # 批量读取所有符合条件的Parquet文件 merged_df = spark.read.parquet(*spark_readable_paths) # coalesce(1)将所有数据合并到单个分区,写入后只会生成1个Parquet数据文件 merged_df.coalesce(1).write.mode("overwrite").parquet(target_path) # --------------------------可选步骤-------------------------- # 如果需要把输出目录下的part-xxx.parquet重命名为固定文件名、清理多余的标记文件,放开以下代码 # local_target = target_path.replace('dbfs:/', '/dbfs/') # for f in os.listdir(local_target): # f_path = os.path.join(local_target, f) # if f.startswith('part-') and f.endswith('.parquet'): # # 重命名为你需要的固定文件名 # os.rename(f_path, os.path.join(local_target, 'merged_result.parquet')) # elif f.startswith('_'): # # 删除_SUCCESS、_committed等标记文件 # os.remove(f_path)
关键说明
coalesce(1)会把全量数据汇聚到单个Executor节点执行写入,仅适合总数据量在100GB以内的场景,数据量过大时容易触发Executor OOM,可根据集群内存情况调整分区数,或先做数据清洗过滤后再合并- 写入模式
mode("overwrite")表示目标路径存在时直接覆盖,可根据业务需求替换为append(追加写入)、error(路径存在则报错,为默认值)、ignore(路径存在则跳过写入) - 若集群配置了S3/ADLS等存储层的文件提交协议,直接重命名part文件的可选步骤可能存在权限问题,无强需求可不执行,只要写入时用
coalesce(1),输出目录下就只会有1个Parquet数据文件,不影响后续读取使用。
内容的提问来源于stack exchange,提问作者Keylie
相关产品推荐
相关产品推荐

