Foundry中如何读取快照视图输出数据集上次构建后新增行?
问题描述
我将输出数据集设置为快照事务视图,尝试读取上次运行后输出数据集中的新增行,但得到0行结果。有两个核心疑问:
- 针对快照视图的输出数据集,是否可以访问其上次运行后新增的行?
- 如果不可行,有没有管道逻辑方案可以实现该需求?
另外,我需要检查同时存在于输入和输出数据集中的行,是否也存在于输出数据集上次运行后新增的行中。
输出数据集的事务视图为快照类型,相关代码如下:
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
相关产品推荐
相关产品推荐

