如何在Foundry中基于Wellname匹配替换Output_df行并增量更新
在Palantir Foundry中实现DataFrame增量更新:替换匹配行并新增行
我有两个DataFrame:input_df和output_df,两者都包含Wellname列。需要实现:
- 将
output_df中与input_df的Wellname匹配的所有行,替换为input_df中对应的行 - 同时添加
input_df中output_df不存在的Wellname对应的新行
示例数据
增量转换前的output_df:
Wellname WellType Platform 0 E17 Producer DWG 1 E17 Producer DWG 2 E17 Producer DWG 3 E20Y Producer DWG 4 E20Y Producer DWG 5 E20Y Producer DWG 6 E20Y Producer DWG 7 E20Y Producer DWG
input_df:
Wellname WellType Platform 0 E17 Producer CH 1 E17 Producer CH 2 E17 Producer CH 3 E21 Producer DWG 4 E21 Producer DWG 5 E21 Producer DWG
期望的增量转换后output_df:
Wellname WellType Platform 0 E17 Producer CH 1 E17 Producer CH 2 E17 Producer CH 3 E20Y Producer DWG 4 E20Y Producer DWG 5 E20Y Producer DWG 6 E20Y Producer DWG 7 E20Y Producer DWG 8 E21 Producer DWG 9 E21 Producer DWG 10 E21 Producer DWG
我尝试的代码(无法实现需求):
from transforms.api import transform, incremental, Input, Output, configure from pyspark.sql import types as T schema = T.StructType([ T.StructField('Wellname', T.StringType()), T.StructField('WellType', T.StringType()), T.StructField('Platform', T.StringType()), ]) @configure(profile=['KUBERNETES_NO_EXECUTORS']) @incremental(require_incremental=True, snapshot_inputs=["input_df"]) @transform( input_df=Input('ri.foundry.lava-catalog.dataset.7ef47ff2-7015-4dad-aec9-aa3075a63a96'), output_df=Output('ri.foundry.lava-catalog.dataset.524cee15-7ac5-410c-96e8-b205bac1cee8') ) def incremental_filter(input_df, output_df): new_df = input_df.dataframe() new_df = new_df.unionByName(output_df.dataframe('previous', schema)) mode = 'modify' new_df = new_df.select('Wellname', 'WellType', 'Platform') output_df.set_mode(mode) output_df.write_dataframe()
解决方案
原代码的问题在于直接用unionByName会保留新旧重复行,没有实现"替换"逻辑。正确思路是先过滤掉历史数据中与输入匹配的行,再合并新数据:
from transforms.api import transform, incremental, Input, Output, configure from pyspark.sql import types as T from pyspark.sql.functions import col schema = T.StructType([ T.StructField('Wellname', T.StringType()), T.StructField('WellType', T.StringType()), T.StructField('Platform', T.StringType()), ]) @configure(profile=['KUBERNETES_NO_EXECUTORS']) @incremental(require_incremental=True, snapshot_inputs=["input_df"]) @transform( input_df=Input('ri.foundry.lava-catalog.dataset.7ef47ff2-7015-4dad-aec9-aa3075a63a96'), output_df=Output('ri.foundry.lava-catalog.dataset.524cee15-7ac5-410c-96e8-b205bac1cee8') ) def incremental_update(input_df, output_df): # 获取输入数据和历史输出数据 input_data = input_df.dataframe() prev_output = output_df.dataframe('previous', schema) # 提取input中存在的Wellname集合 existing_wells = input_data.select(col('Wellname')).distinct() # 过滤历史数据:保留input中没有的Wellname对应的行 filtered_prev = prev_output.join(existing_wells, on='Wellname', how='anti') # 合并过滤后的历史数据和输入数据 final_df = filtered_prev.unionByName(input_data) # 设置modify模式并写入,替换旧数据 output_df.set_mode('modify') output_df.write_dataframe(final_df)
关键说明
- Anti Join:用
how='anti'的join操作,高效过滤掉历史数据中与输入数据Wellname匹配的行,实现"替换"的核心逻辑 - UnionByName:确保列名匹配的情况下合并数据,保留所有需要的行
- Modify模式:Foundry的
modify模式会用新数据替换输出数据集的内容,确保最终结果是过滤后的历史数据+新输入数据的组合
内容的提问来源于stack exchange,提问作者Mahammad Ojagzada
相关产品推荐
相关产品推荐

