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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.02 13:05:16