如何在Foundry代码仓库中读取output_df的全部数据?
问题描述
需要访问output_df的全部数据以与input_df进行对比。经查看,input_df和output_df的事务视图均为snapshot,当前output_df有250行,但使用"current"、"added"和"previous"尝试读取其全部数据时,在预览模式下均返回0行。相关代码如下:
@configure(profile=['KUBERNETES_NO_EXECUTORS']) @incremental(semantic_version=17) @transform( input_df=Input('ri.foundry.lava-catalog.dataset.c855c91b-3e73-4803-84bd-7d35c45f724c'), output_df=Output('ri.foundry.lava-catalog.dataset.4802782e-4436-4bdf-87f3-5457245574c1') ) def incremental_filter(input_df, output_df): df_new = input_df.dataframe('added') df_new = df_new.withColumn('Start', to_timestamp(col('Start'), 'yyyy-MM-dd HH:mm:ss')) df_new = df_new.withColumn('End', to_timestamp(col('End'), 'yyyy-MM-dd HH:mm:ss')) print("df_new columns are {}".format(df_new.columns)) print('----------------------------------') # Load previous dataframe print(output_df.dataframe('current', schema=schema).localCheckpoint().count()) df_previous = output_df.dataframe('current', schema=schema) print("df_previous current count: ", df_previous.count()) #0 rows df_previous = output_df.dataframe('previous', schema=schema) print("df_previous previous count: ", df_previous.count()) #0 rows df_previous = output_df.dataframe('added', schema=schema) print("df_previous added count: ", df_previous.count()) #0 rows #Doing some comparisons here # ................... # -------------------------- mode = 'replace' output_df.set_mode(mode) # Write the output dataframe output_df.write_dataframe(df_union)
请问如何在Foundry代码仓库中获取output_df的全部数据?
解决方案
1. 理解预览模式的限制
预览模式下,增量函数的Output对象无法返回历史数据,current/previous/added都会返回空数据集——这是因为预览是模拟增量运行流程,不会加载真实的历史输出状态。要读取真实全量数据,需避开预览模式的限制。
2. 将输出数据集作为独立Input引入
如果需要在函数内获取output_df的全量历史数据,不要通过Output对象的dataframe()方法读取,而是将目标输出数据集作为额外的Input参数引入,直接读取全量:
from transforms.api import Input, Output, transform, incremental, configure from pyspark.sql.functions import col, to_timestamp @configure(profile=['KUBERNETES_NO_EXECUTORS']) @incremental(semantic_version=17) @transform( input_df=Input('ri.foundry.lava-catalog.dataset.c855c91b-3e73-4803-84bd-7d35c45f724c'), output_df=Output('ri.foundry.lava-catalog.dataset.4802782e-4436-4bdf-87f3-5457245574c1'), # 新增:将输出数据集作为独立Input引入 output_full=Input('ri.foundry.lava-catalog.dataset.4802782e-4436-4bdf-87f3-5457245574c1') ) def incremental_filter(input_df, output_df, output_full): df_new = input_df.dataframe('added') df_new = df_new.withColumn('Start', to_timestamp(col('Start'), 'yyyy-MM-dd HH:mm:ss')) df_new = df_new.withColumn('End', to_timestamp(col('End'), 'yyyy-MM-dd HH:mm:ss')) print("df_new columns are {}".format(df_new.columns)) print('----------------------------------') # 读取输出数据集的全量数据 df_output_full = output_full.dataframe() print("全量输出数据行数: ", df_output_full.count()) # 此时会返回真实的250行 # 后续对比逻辑... # ................... # -------------------------- mode = 'replace' output_df.set_mode(mode) # Write the output dataframe output_df.write_dataframe(df_union)
3. 其他可行方案
- 全量运行函数:在函数运行时选择「全量运行」而非增量运行,此时
output_df.dataframe('current')会加载输出数据集的历史全量数据(仅非预览模式下有效)。 - 直接查看数据集:在Foundry数据集页面直接预览
output_df的内容,确认数据存在后再回到代码中做逻辑开发。
内容的提问来源于stack exchange,提问作者Mahammad Ojagzada
相关产品推荐
相关产品推荐

