Xây dựng Data Pipeline với Airflow
Volker Janz
Senior Developer Advocate at Astronomer
{{ ds }} gắn mỗi lần chạy với một ngày, phù hợp cho Dags lập lịch theo thời giands là None, tức là không có sẵnpartition_key từ upstream xuống downstream
from airflow.sdk import dag, CronPartitionTimetable
@dag(schedule=CronPartitionTimetable("0 0 * * *", timezone="UTC"))
def sales_pipeline():
...
partition_keyTrong tác vụ Python:
@task
def process_data(dag_run=None):
partition_key = dag_run.partition_key
print(f"Processing: {partition_key}")
Trong mẫu SQL:
DELETE FROM daily_summary
WHERE order_date = '{{ dag_run.partition_key[:10] }}';
from airflow.sdk import dag, task, Asset, CronPartitionTimetable
daily_sales = Asset("daily_sales")
@dag(schedule=CronPartitionTimetable("0 0 * * *", timezone="UTC"))
def sales_pipeline():
@task(outlets=[daily_sales])
def load_data(**context):
...
from airflow.sdk import dag, task, Asset, PartitionedAssetTimetable
daily_sales = Asset("daily_sales")
@dag(schedule=PartitionedAssetTimetable(assets=daily_sales))
def sales_report():
@task
def generate_report(dag_run=None):
partition_key = dag_run.partition_key
print(f"Report for: {partition_key}")
from airflow.sdk import (
dag, PartitionedAssetTimetable,
StartOfDayMapper,
)
@dag(schedule=PartitionedAssetTimetable(
assets=daily_sales,
partition_mapper_config={
daily_sales: StartOfDayMapper()
},
))
def sales_report():
...
Trước khi mapping:
2026-04-23T00:00:00
Sau StartOfDayMapper:
2026-04-23
$$
StartOfHourMapper, StartOfWeekMapper, StartOfMonthMapperAllowedKeyMapper cho vùng, phòng ban
CronPartitionTimetable gắn partition key vào mỗi lần chạyStartOfDayMapper chuẩn hóa timestamp thành ngàyPartitionedAssetTimetable kích hoạt và kế thừa key đã mapXây dựng Data Pipeline với Airflow