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

如何实现首次全量刷新后持续增量追加的混合同步模式?

自定义Report源连接器同步模式与归一化问题

当前配置

Source: Custom Report API (custom connector)
Sync Requirements:
First sync: Full Refresh | Overwrite (拉取全部历史数据)
Subsequent syncs: Append only (仅拉取上一财年数据)
Using Basic Normalization
Tracking sync state in external database (sync_completed_date)

示例代码

class Report(ReportStream):
    report_name = "Report"
    primary_key = ""
    
    def path(self, ...):
        if not self.sync_completed_date:
            # First sync - get all historical data
            start_year = date.today().year - 30
            start_date = f"{start_year}-01-01"
        else:
            # Subsequent syncs - get only previous fiscal year data
            start_date, end_date = self.get_fiscal_year_dates(self.fiscal_month)
        return (
            f"/reports/{self.report_name}&"
            f"start_date={start_date}&"
            f"end_date={end_date}&"
            f"columns={self.output_columns}&"
        )

遇到的问题

  • 首次同步设置为「Full Refresh | Overwrite」时正常,能拉取全部历史数据
  • 第二次同步时,连接器正确拉取了上一财年的数据,但归一化仍执行覆盖操作,清除了所有历史数据,导致丢失上一财年之前的数据

疑问

  1. 如何控制归一化行为,使其在首次同步后执行追加操作,即使UI设置为「Full Refresh | Overwrite」?
  2. 是否有推荐的模式来实现这种混合同步行为(首次全量刷新,之后始终追加)?

解决方案

1. 动态控制同步模式与归一化行为

不要依赖UI的固定设置,直接在连接器代码中根据sync_completed_date状态动态返回对应的同步模式和目标同步模式:

from airbyte_cdk.models import SyncMode, DestinationSyncMode

class Report(ReportStream):
    report_name = "Report"
    primary_key = ""
    
    def path(self, ...):
        # 原path逻辑保留
        if not self.sync_completed_date:
            start_year = date.today().year - 30
            start_date = f"{start_year}-01-01"
        else:
            start_date, end_date = self.get_fiscal_year_dates(self.fiscal_month)
        return (
            f"/reports/{self.report_name}&"
            f"start_date={start_date}&"
            f"end_date={end_date}&"
            f"columns={self.output_columns}&"
        )
    
    def get_sync_mode(self, context):
        # 首次同步用全量,后续用增量
        if not self.sync_completed_date:
            return SyncMode.FULL_REFRESH
        else:
            return SyncMode.INCREMENTAL
    
    def get_destination_sync_mode(self, context):
        # 首次同步覆盖,后续追加
        if not self.sync_completed_date:
            return DestinationSyncMode.OVERWRITE
        else:
            return DestinationSyncMode.APPEND

这样连接器会自动根据状态切换行为,忽略UI的初始设置,确保首次同步覆盖全量数据,后续只追加增量数据。

2. 推荐的混合同步实现模式

  • 状态驱动判断:始终以外部存储的sync_completed_date作为核心判断依据,不要依赖UI或临时变量
  • 状态更新保障:每次同步完成后,务必将sync_completed_date更新为本次同步的结束日期(或当前时间),确保后续同步能正确进入增量模式
  • 边界容错处理:如果sync_completed_date丢失或异常,可额外检查目标表是否存在数据——若表为空则执行全量覆盖,否则直接执行增量追加,避免误删数据
  • 主键约束:如果报表数据有唯一主键,建议设置primary_key,这样即使后续追加重复数据,归一化阶段也能自动处理(避免重复插入)

内容的提问来源于stack exchange,提问作者frank

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.12 11:13:17