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

Foundry中如何读取快照视图输出数据集上次构建后新增行?

问题描述

我将输出数据集设置为快照事务视图,尝试读取上次运行后输出数据集中的新增行,但得到0行结果。有两个核心疑问:

  1. 针对快照视图的输出数据集,是否可以访问其上次运行后新增的行?
  2. 如果不可行,有没有管道逻辑方案可以实现该需求?
    另外,我需要检查同时存在于输入和输出数据集中的行,是否也存在于输出数据集上次运行后新增的行中。

输出数据集的事务视图为快照类型,相关代码如下:

from transforms.api import transform, incremental, Input, Output, configure
from pyspark.sql import types as T
from pyspark.sql import functions as F

schema = T.StructType(
    [
        T.StructField("wellname", T.StringType()),
        T.StructField("start_date", T.TimestampType()),
        T.StructField("woe_limit_psia", T.DoubleType()),
        T.StructField("end_date", T.TimestampType()),
    ]
)


@configure(profile=["KUBERNETES_NO_EXECUTORS"])
@incremental(
    # require_incremental=True, snapshot_inputs=["input_df"], semantic_version=13
    snapshot_inputs=["input_df", "input_df2"], require_incremental=True
)
@transform(
    input_df=Input(
        "ri.foundry.lava-catalog.dataset.974e09fa-fba6-4867-a7b1-f860ab7a7046"
    ),
    output_df=Output(
        "ri.foundry.lava-catalog.dataset.c3407f67-2798-4049-8455-586a025a0b65"
    )
)

def incremental_filter(input_df, output_df):
    df_current = input_df.dataframe("added")
    df_history = output_df.dataframe('previous', schema=schema)
    new_output = output_df.dataframe('added', schema=schema)
    
    # 需检查输入和输出都存在的行是否在输出新增行中

    previous_rows = df_current.join(df_history, on=["wellname", "woe_limit_psia"], how="inner")
    repeated_rows = previous_rows.join(new_output, on=['wellname', 'woe_limit_psia'], how='left_anti')
    # repeated_rows包含全部previous_rows数据,推测因new_output获取到0行
解决方案与说明

1. 快照视图能否访问上次运行新增行?

不能。快照事务视图的核心特性是每次运行都会生成完整的数据集快照,不会保留增量变更的追踪信息。当调用output_df.dataframe('added', schema=schema)时,快照视图没有"新增行"的概念,因此返回0行是预期行为。

2. 实现需求的管道逻辑方案

要追踪输出数据集的增量变更,可通过以下两种方案实现:

方案一:新增增量式中间数据集

  • 创建一个新的增量事务视图数据集作为中间层,专门用于记录每次运行的增量变更。
  • 在当前管道中,将原本输出到快照视图的数据同时写入这个增量数据集。
  • 后续需要读取上次运行新增行时,直接从该增量数据集调用.dataframe('added')获取。

方案二:在现有快照数据集中维护增量标记

  • 扩展输出数据的Schema,新增一个_load_timestamp字段,每次运行时用当前时间戳标记新写入的行。
  • 读取上次运行新增行时,通过筛选_load_timestamp大于上次运行时间戳的行来获取增量数据。需额外维护上次运行的时间戳(可存储在一个小型辅助数据集中)。

3. 检查重复行的逻辑调整

结合上述方案,调整代码实现"输入和输出都存在的行是否在输出新增行中"的检查需求:

基于增量中间数据集的代码示例

from transforms.api import transform, incremental, Input, Output, configure
from pyspark.sql import types as T
from pyspark.sql import functions as F

schema = T.StructType(
    [
        T.StructField("wellname", T.StringType()),
        T.StructField("start_date", T.TimestampType()),
        T.StructField("woe_limit_psia", T.DoubleType()),
        T.StructField("end_date", T.TimestampType()),
    ]
)

@configure(profile=["KUBERNETES_NO_EXECUTORS"])
@incremental(
    snapshot_inputs=["input_df"], require_incremental=True
)
@transform(
    input_df=Input("ri.foundry.lava-catalog.dataset.974e09fa-fba6-4867-a7b1-f860ab7a7046"),
    output_df=Output("ri.foundry.lava-catalog.dataset.c3407f67-2798-4049-8455-586a025a0b65"),
    # 新增增量中间数据集,事务类型设为增量视图
    incremental_output=Output("ri.foundry.lava-catalog.dataset.xxx-xxx-xxx-xxx-xxx")
)
def incremental_filter(input_df, output_df, incremental_output):
    df_current = input_df.dataframe("added")
    df_history = output_df.dataframe('previous', schema=schema)
    
    # 从增量中间数据集获取上次运行新增的行
    last_run_added = incremental_output.dataframe('added', schema=schema)
    
    # 检查输入与输出都存在的行是否不在上次新增行中
    previous_rows = df_current.join(df_history, on=["wellname", "woe_limit_psia"], how="inner")
    repeated_rows = previous_rows.join(last_run_added, on=['wellname', 'woe_limit_psia'], how='left_anti')
    
    # 将本次结果写入快照输出和增量中间数据集
    output_df.write_dataframe(repeated_rows)
    incremental_output.write_dataframe(repeated_rows)

基于时间戳标记的代码示例

from transforms.api import transform, incremental, Input, Output, configure
from pyspark.sql import types as T
from pyspark.sql import functions as F

# 扩展Schema,新增加载时间戳字段
extended_schema = T.StructType(
    [
        T.StructField("wellname", T.StringType()),
        T.StructField("start_date", T.TimestampType()),
        T.StructField("woe_limit_psia", T.DoubleType()),
        T.StructField("end_date", T.TimestampType()),
        T.StructField("_load_timestamp", T.TimestampType()),
    ]
)

@configure(profile=["KUBERNETES_NO_EXECUTORS"])
@incremental(
    snapshot_inputs=["input_df"], require_incremental=True
)
@transform(
    input_df=Input("ri.foundry.lava-catalog.dataset.974e09fa-fba6-4867-a7b1-f860ab7a7046"),
    output_df=Output("ri.foundry.lava-catalog.dataset.c3407f67-2798-4049-8455-586a025a0b65"),
    # 存储上次运行时间戳的辅助数据集
    last_run_ts=Output("ri.foundry.lava-catalog.dataset.yyy-yyy-yyy-yyy-yyy")
)
def incremental_filter(input_df, output_df, last_run_ts):
    df_current = input_df.dataframe("added")
    df_history = output_df.dataframe('previous', schema=extended_schema)
    
    # 获取上次运行的时间戳
    last_ts = last_run_ts.dataframe().select(F.max("_load_timestamp")).first()[0] if last_run_ts.exists() else None
    current_ts = F.current_timestamp()
    
    # 筛选上次运行后新增的行
    last_run_added = df_history.filter(F.col("_load_timestamp") > last_ts) if last_ts else df_history.limit(0)
    
    # 检查重复行逻辑(移除时间戳字段后关联)
    previous_rows = df_current.join(df_history.drop("_load_timestamp"), on=["wellname", "woe_limit_psia"], how="inner")
    repeated_rows = previous_rows.join(last_run_added.drop("_load_timestamp"), on=['wellname', 'woe_limit_psia'], how='left_anti')
    
    # 标记本次加载时间戳并写入快照数据集
    output_data = repeated_rows.withColumn("_load_timestamp", current_ts)
    output_df.write_dataframe(output_data)
    
    # 更新辅助数据集中的上次运行时间戳
    last_run_ts.write_dataframe(spark.createDataFrame([(current_ts,)], schema=T.StructType([T.StructField("_load_timestamp", T.TimestampType())])))

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.18 11:34:53