Airflow によるデータパイプラインの構築
Volker Janz
Senior Developer Advocate at Astronomer
@task.bash
def export_data():
return "cp /data/sales.csv /export/sales_{{ ds }}.csv"
{{ ds }} は論理日付を YYYY-MM-DD 形式でレンダリングしますbash_command や sql などのテンプレート対応フィールドでのみ使用可能です
実行 1(3月31日):
実行 2(3月31日を再実行):

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, )
$$
UPSERT や MERGE による方法も使用できます


Airflow によるデータパイプラインの構築