在Foundry增量转换中如何读取输出数据集的完整历史数据?
问题
我将输入数据集(代码工作簿Python转换生成的Snapshot类型,共86行)处理后,以Append事务类型将输出数据集(共123行)上传至Foundry,格式为Parquet。现在需要读取该输出数据集的全部123行数据,与输入数据集对比后,用新数据替换整个输出数据集。
测试时我尝试读取输出数据集的历史版本并与输入数据集做Union,但最终输出仅得到86行(和输入数据一致),推测未成功读取到输出数据集的完整数据。请问如何读取输出数据集的全部123行数据?
测试代码
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"] ) @transform( input_df=Input( "ri.foundry.lava-catalog.dataset.974e09fa-fba6-4867-a7b1-f860ab7a7046" ), output_df=Output( "ri.foundry.lava-catalog.dataset.95804158-098a-4f0c-afea-64abd10fbc1f" ) ) def incremental_filter(input_df, output_df): df_current = input_df.dataframe("added") columns_in_order = [field.name for field in schema.fields] df_current = df_current.select(columns_in_order) df_union = df_current.unionByName(output_df.dataframe(mode="previous", schema=schema)) mode = "replace" df_union.localCheckpoint(eager=True) output_df.set_mode(mode) # Write the output dataframe output_df.write_dataframe(df_union)
测试结果
仅得到86行数据(与输入数据集行数一致)
解决方法
问题出在读取输出数据集的模式上:你使用了mode="previous",该模式在增量转换中只会读取由当前增量转换生成的上一个版本的数据,但你的输出数据集是通过Append手动上传的,并非当前增量转换生成,因此无法读取到这部分历史数据,导致Union后只有输入的86行。
要读取输出数据集的全部历史数据(包括之前Append上传的123行),需要将读取模式改为mode="full",该模式会加载输出数据集的所有完整数据。
修改后的完整代码
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"] ) @transform( input_df=Input( "ri.foundry.lava-catalog.dataset.974e09fa-fba6-4867-a7b1-f860ab7a7046" ), output_df=Output( "ri.foundry.lava-catalog.dataset.95804158-098a-4f0c-afea-64abd10fbc1f" ) ) def incremental_filter(input_df, output_df): df_current = input_df.dataframe("added") columns_in_order = [field.name for field in schema.fields] df_current = df_current.select(columns_in_order) # 改为full模式读取输出数据集全部数据 df_union = df_current.unionByName(output_df.dataframe(mode="full", schema=schema)) mode = "replace" df_union.localCheckpoint(eager=True) output_df.set_mode(mode) # Write the output dataframe output_df.write_dataframe(df_union)
如果你的最终需求是替换整个输出数据集,也可以考虑移除@incremental装饰器,直接以全量模式运行转换,逻辑会更简洁:读取输入全量数据和输出全量数据,对比处理后直接替换输出。
内容的提问来源于stack exchange,提问作者Mahammad Ojagzada
相关产品推荐
相关产品推荐

