如何将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
相关产品推荐
相关产品推荐

