模板化、幂等性与回补

使用 Airflow 构建数据流水线

Volker Janz

Senior Developer Advocate at Astronomer

Airflow 中的 Jinja 模板

@task.bash
def export_data():
    return "cp /data/sales.csv /export/sales_{{ ds }}.csv"

 

  • {{ ds }} 会渲染为 YYYY-MM-DD 格式的逻辑日期
  • 每次运行都有自己的日期,因此同一任务会生成每次运行不同的输出
  • 仅适用于可模板化字段(如 bash_commandsql
使用 Airflow 构建数据流水线

幂等性难题

 

运行 1(3 月 31 日):

  • 插入 3 月 31 日的销售数据
  • 结果:3 行

运行 2(重跑 3 月 31 日):

  • 再次插入 3 月 31 日的销售数据
  • 结果:6 行(重复!)

重跑导致重复行

使用 Airflow 构建数据流水线

先删后插模式

staged_rows = SQLExecuteQueryOperator(
    task_id="get_staged_rows",
    conn_id="duckdb_default",
    sql="SELECT * FROM staging WHERE date = '{{ ds }}'",
)

load = SQLInsertRowsOperator( task_id="load_sales", conn_id="duckdb_default", table_name="sales", columns=["date", "product", "amount"], preoperator="DELETE FROM sales WHERE date = '{{ ds }}';", rows=staged_rows.output, )

$$

  • 也可使用 UPSERTMERGE 方法

先删后插

使用 Airflow 构建数据流水线

回补(Backfilling)

 

  • 重新处理一段历史日期
  • Airflow 为每个逻辑日期创建一个 Dag 运行
  • 结合幂等性,回补是安全

回补取决于调度

使用 Airflow 构建数据流水线

在界面中回补

  • 触发一个 Dag,选择 Backfill,设置 FromTo 日期
  • 可选择触发 Missing RunsMissing and Errored RunsAll Runs
  • 可用 Max Active Runs 设置并发,并调整执行顺序
  • 点击 Run Backfill 后开始流程并创建运行

Airflow 回补界面

使用 Airflow 构建数据流水线

让我们一起练习吧!

使用 Airflow 构建数据流水线

Preparing Video For Download...