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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.03 04:48:29