如何在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}")
关键说明
- 缓存基准文件:
cache()方法将df_mirror驻留内存,大幅提升循环处理效率 - 循环遍历文件:直接迭代
sorted后的test列表,逐个获取文件路径 - 临时视图复用:每次循环覆盖
deltas视图,确保SQL始终读取当前文件的数据 - 结果处理:可根据需求替换保存逻辑,如写入数据库、合并结果集等
内容的提问来源于stack exchange,提问作者nox8315
相关产品推荐
相关产品推荐

