基于增量运行能力调整的PySpark历史记录构建技术咨询
问题背景与解决方案
我是编程新手,正在自学Python和PySpark,需要基于每日变更构建历史数据集。需要定期升级语义版本(SEMANTIC_VERSION),但不想丢失已收集的历史数据。期望实现:当任务可增量运行时执行常规增量转换;当任务无法增量运行时,将当前快照数据与已收集的历史数据合并。以下是我的尝试代码及基于反馈优化后的可用代码:
初始尝试代码
SEMANTIC_VERSION = 1 # 当任务无法增量运行时 # 将当前快照数据与已收集的历史数据合并 if cannot_not_run_incrementally: @transform( history=Output(historical_output), backup=Input(historical_output_backup), source=Input(order_input), ) def my_compute_function(source, history, backup, ctx): input_df = ( source.dataframe() .withColumn('record_date', F.current_date()) ) old_df = backup.dataframe() joined = old_df.unionByName(input_df) joined = joined.distinct() history.write_dataframe(joined) # 当任务可增量运行时,执行常规增量转换 else: @incremental(snapshot_inputs=['source'], semantic_version=SEMANTIC_VERSION) @transform( history=Output(historical_output), backup=Output(historical_output_backup), source=Input(order_input), ) def my_compute_function(source, history, backup): input_df = ( source.dataframe() .withColumn('record_date', F.current_date()) ) history.write_dataframe(input_df.distinct() .subtract(history.dataframe('previous', schema=input_df.schema))) backup.set_mode("replace") backup.write_dataframe(history.dataframe())
优化后的可用代码
SEMANTIC_VERSION = 3 @incremental(snapshot_inputs=['source'], semantic_version=SEMANTIC_VERSION) @transform( history=Output(), backup=Output(), source=Input(), ) def compute(ctx, history, backup, source): # 增量运行模式 if ctx.is_incremental: input_df = ( source.dataframe() .withColumn('record_date', F.current_date()) ) history.write_dataframe(input_df.subtract(history.dataframe('previous', schema=input_df.schema))) backup.set_mode("replace") backup.write_dataframe(history.dataframe().distinct()) # 非增量运行模式 else: input_df = ( source.dataframe() .withColumn('record_date', F.current_date()) ) backup.set_mode('modify') # 如果想要重新开始可以使用replace模式 backup.write_dataframe(input_df) history.set_mode('replace') history.write_dataframe(backup.dataframe().distinct())
内容的提问来源于stack exchange,提问作者eruhl06
相关产品推荐
相关产品推荐

