airflow-dag-patternslisted
Install: claude install-skill sandbaseai/workbuddy-skill
# Airflow DAG Patterns
设计可观察、可重试、可回放的数据编排。重点是 DAG 的时间语义、幂等任务、依赖等待、失败处理、测试和部署边界,而不是替用户直接触发生产任务。
## 使用边界
- 开始前确认 Airflow 版本、provider、时区、DAG 提交、数据范围、owner、SLA、连接引用和目标环境。
- 默认只读检查 DAG、配置、任务日志、依赖和测试;验证使用本地/CI/隔离环境及合成或脱敏数据。
- 不擅自触发生产 DAG、重试、清理、`backfill`、`clear`、`pause/unpause`、连接/变量修改或数据写入;这些操作需要明确授权。
- 日志、XCom、样本和报告不得保存 secret、token、完整客户行或可还原的个人标识。
## DAG 设计
每个 DAG 说明业务目的、输入分区、输出、粒度、调度时区、数据就绪条件、最大并发、超时、SLA、重跑和所有者。任务应满足:
- 幂等:同一逻辑日期重复执行不会重复写入或产生不同结果;用分区键/幂等键和事务边界约束写入。
- 原子:失败不会留下“成功”标记或半成品;临时结果在校验后再发布。
- 可增量:按逻辑日期和 watermark 处理,明确迟到数据、删除和更新语义。
- 可观察:结构化日志、耗时/行数/新鲜度指标、失败上下文和告警都能追溯到 DAG run 与 task instance。
使用 TaskFlow API 或清晰的 Operator 封装任务逻辑。DAG 文件保持轻量,避免导入时联网、查询数据库、读取不稳定的全局状态或执行重计算。
```python
from datetime import datetime, timedelta
from airflow.decorators import dag, task
@dag(
dag_id="daily_orders",
schedule="0 6 * * *",
start_date=datetime(2024, 1, 1),
catchup=False,
max_active_runs=1,
dagrun_timeout=timedelta(hours=2),
tags=["etl"],
)
def daily_orders():
@task(retries=3, retry_exponential_backoff=True,
execution_timeout=timedelta(minutes=30))
def extract(logical_date=None):
# Use the logical date/partition, not wall-clock "now".
return {"partition": logical_date.strftime("%Y-%m-%d")}
@task
def validate(input_ref):
# Validate completeness and schema before publishing output.
return input_ref
validate(extract())
daily_orders()
```
## 依赖、传感器与分支
用 `upstream >> downstream` 表达依赖;