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

基于增量运行能力调整的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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.24 22:48:51