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

如何将PySpark DataFrame每一行分别写入独立文本文件

PySpark 逐行生成独立文本文件实现方案

注意:5万行数据量级下禁止直接调用collect()全量拉取数据到Driver节点,若Content字段单条内容过长极易触发Driver OOM,根据数据规模可选择以下两种实现:


方案1:小数据量场景(总内容体积<1GB):Driver端遍历写入

适合单条Content长度短、总数据量不大的场景,直接通过遍历逐行调用文件写入接口,逻辑简单易调试。

# 替换为你的目标存储路径,Databricks环境可直接用DBFS路径
OUTPUT_DIR = "/dbfs/FileStore/tables/individual_texts/"
# 提前创建输出目录
dbutils.fs.mkdirs(OUTPUT_DIR)

# 用toLocalIterator替代collect,流式逐批拉取数据到Driver,避免内存溢出
for row in df_test.toLocalIterator():
    # 提取文件名和内容,可按需对ID做特殊字符替换,避免非法文件名报错
    file_name = f"{row['ID'].replace('/', '_').replace(':', '_')}.txt"
    file_path = f"{OUTPUT_DIR}{file_name}"
    # 写入文件,overwrite=True表示重名文件直接覆盖
    dbutils.fs.put(file_path, row["Content"], overwrite=True)

该方案特点:

  • 代码量少,逻辑直观,调试成本低
  • 所有写入操作集中在Driver节点串行执行,数据量大时写入速度慢

方案2:大规模数据场景:Executor端分布式并行写入

5万行数据优先选这个方案,利用Spark分布式计算能力在各Executor节点并行写入,速度是方案1的数倍到数十倍,无Driver内存溢出风险。

import os

# 替换为你的目标存储路径
OUTPUT_DIR = "/dbfs/FileStore/tables/individual_texts/"
dbutils.fs.mkdirs(OUTPUT_DIR)

def write_rows_to_files(partition_iter):
    # 每个数据分区内独立执行写入逻辑,无需跨节点传输数据
    for row in partition_iter:
        # 过滤ID中的非法文件名字符
        safe_id = row["ID"].replace("/", "_").replace("\\", "_").replace(":", "_")
        file_path = os.path.join(OUTPUT_DIR, f"{safe_id}.txt")
        # 用原生Python IO接口写入,兼容性比dbutils更好,支持大文件写入
        with open(file_path, "w", encoding="utf-8") as f:
            f.write(row["Content"])
    # 返回空迭代器满足mapPartitions的接口要求
    return iter([])

# 按需调整分区数,5万行数据设置8-16个分区即可获得最优并行度
df_test.repartition(16).rdd.mapPartitions(write_rows_to_files).count()

该方案注意事项:

  • Databricks Runtime默认支持/dbfs前缀的路径本地挂载,若使用其他Spark发行版,可替换为本地磁盘路径或HDFS路径
  • 写入前建议先对ID字段去重,避免相同ID的文件被后续写入的内容覆盖
  • 若存储路径为对象存储(比如S3、OSS),需要提前在集群配置好对应访问密钥

结果校验

写入完成后可通过以下命令校验文件是否符合预期:

# 列出输出目录下所有生成的文件
file_list = dbutils.fs.ls(OUTPUT_DIR)
print(f"共生成{len(file_list)}个文本文件")

# 读取单个文件验证内容
print(dbutils.fs.head(f"{OUTPUT_DIR}A1234.txt"))

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.27 20:09:19