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

如何在PySpark中循环处理CSV文件并执行全外连接

问题:如何循环处理找到的CSV文件并与基准文件执行全外连接?

已通过以下代码找到并排序了不同文件夹下的5个目标CSV文件:

list_files = glob.glob("/t/main_folder/*/file_*[0-9].csv")
test = sorted(list_files, key = lambda x:x[-5:])

当前单文件处理代码如下,但需要实现循环逐个读取上述文件,与预加载的/t/org_file.csv执行全外连接:

df_deltas = spark.read.format("csv").schema(schema).option("header","true")\
.option("delimiter",";").load(test)
df_mirror = spark.read.format("csv").schema(schema).option("header","true")\
.option("delimiter",",").load("/t/org_file.csv").cache()
df_deltas.createOrReplaceTempView("deltas")
df_mirror.createOrReplaceTempView("mirror")
df_mir2=spark.sql("""select 
coalesce (deltas.DATA_ACTUAL_DATE,mirror.DATA_ACTUAL_DATE) as DATA_ACTUAL_DATE,
coalesce (deltas.DATA_ACTUAL_END_DATE,mirror.DATA_ACTUAL_END_DATE) as DATA_ACTUAL_END_DATE,
coalesce (deltas.ACCOUNT_RK,mirror.ACCOUNT_RK) as ACCOUNT_RK,
coalesce (deltas.ACCOUNT_NUMBER,mirror.ACCOUNT_NUMBER) as ACCOUNT_NUMBER,
coalesce (deltas.CHAR_TYPE,mirror.CHAR_TYPE) as CHAR_TYPE,
coalesce (deltas.CURRENCY_RK,mirror.CURRENCY_RK) as CURRENCY_RK,
coalesce (deltas.CURRENCY_CODE,mirror.CURRENCY_CODE) as CURRENCY_CODE,
coalesce (deltas.CLIENT_ID,mirror.CLIENT_ID) as CLIENT_ID,
coalesce (deltas.BRANCH_ID,mirror.BRANCH_ID) as BRANCH_ID,
coalesce (deltas.OPEN_IN_INTERNET,mirror.OPEN_IN_INTERNET) as OPEN_IN_INTERNET
from mirror
full outer join deltas on
deltas.ACCOUNT_RK=mirror.ACCOUNT_RK
""")

解决方案

核心思路

  • 基准文件org_file.csv只需预加载一次并缓存,避免重复IO开销
  • 遍历排序后的文件列表,逐个读取目标CSV文件
  • 每次循环更新临时视图,执行全外连接后处理结果

完整代码示例

# 预加载基准文件并缓存(仅执行一次)
df_mirror = spark.read.format("csv").schema(schema).option("header","true")\
.option("delimiter",",").load("/t/org_file.csv").cache()
df_mirror.createOrReplaceTempView("mirror")

# 遍历每个目标文件路径
for file_path in test:
    # 读取当前目标CSV文件
    df_deltas = spark.read.format("csv").schema(schema).option("header","true")\
    .option("delimiter",";").load(file_path)
    
    # 创建/覆盖临时视图,供SQL调用
    df_deltas.createOrReplaceTempView("deltas")
    
    # 执行全外连接SQL
    df_mir2 = spark.sql("""
        select 
            coalesce(deltas.DATA_ACTUAL_DATE, mirror.DATA_ACTUAL_DATE) as DATA_ACTUAL_DATE,
            coalesce(deltas.DATA_ACTUAL_END_DATE, mirror.DATA_ACTUAL_END_DATE) as DATA_ACTUAL_END_DATE,
            coalesce(deltas.ACCOUNT_RK, mirror.ACCOUNT_RK) as ACCOUNT_RK,
            coalesce(deltas.ACCOUNT_NUMBER, mirror.ACCOUNT_NUMBER) as ACCOUNT_NUMBER,
            coalesce(deltas.CHAR_TYPE, mirror.CHAR_TYPE) as CHAR_TYPE,
            coalesce(deltas.CURRENCY_RK, mirror.CURRENCY_RK) as CURRENCY_RK,
            coalesce(deltas.CURRENCY_CODE, mirror.CURRENCY_CODE) as CURRENCY_CODE,
            coalesce(deltas.CLIENT_ID, mirror.CLIENT_ID) as CLIENT_ID,
            coalesce(deltas.BRANCH_ID, mirror.BRANCH_ID) as BRANCH_ID,
            coalesce(deltas.OPEN_IN_INTERNET, mirror.OPEN_IN_INTERNET) as OPEN_IN_INTERNET
        from mirror
        full outer join deltas on deltas.ACCOUNT_RK = mirror.ACCOUNT_RK
    """)
    
    # 自定义结果处理逻辑:示例为保存到输出目录,文件名保留原文件标识
    output_filename = f"result_{file_path.split('/')[-1]}"
    df_mir2.write.format("csv").option("header","true").save(f"/t/output/{output_filename}")

关键说明

  1. 缓存基准文件:cache()方法将df_mirror驻留内存,大幅提升循环处理效率
  2. 循环遍历文件:直接迭代sorted后的test列表,逐个获取文件路径
  3. 临时视图复用:每次循环覆盖deltas视图,确保SQL始终读取当前文件的数据
  4. 结果处理:可根据需求替换保存逻辑,如写入数据库、合并结果集等

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.31 10:55:50