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

如何编写Airflow自定义DAG装饰器实现多锁互斥机制

自定义DAG装饰器实现自动锁注入

核心思路

通过封装Airflow原生@dag装饰器,自动为目标DAG注入get_lock传感器和free_locks清理任务,同时自动维护任务依赖关系,无需手动编写重复的加锁解锁逻辑。每个DAG可通过参数指定自身需要的锁列表,满足多锁持有的互斥需求。

具体实现步骤

  • 封装原生装饰器:创建自定义装饰器@dag_with_locks,接收原生@dag的所有参数,新增locks参数用于指定当前DAG的锁列表。
  • 自动注入锁任务:在装饰器内部,先生成基础DAG对象,再添加get_lock传感器和free_locks清理任务(标记为teardown确保必执行)。
  • 维护依赖关系:自动识别DAG中的根任务(无上游依赖)和叶子任务(无下游依赖),将get_lock设为所有根任务的上游,free_locks设为所有叶子任务的下游。
  • 边界情况处理:兼容单任务DAG、分支DAG等场景,确保依赖关系正确。

代码实现示例

自定义装饰器代码

from airflow.decorators import dag as original_dag
from airflow.operators.python import PythonOperator
from typing import List, Optional, Callable, Any

# 复用你已实现的create_redis_object、cache、get_lock、free_locks函数
# ...(此处保留你原有的Redis锁相关实现代码)

def dag_with_locks(
    locks: Optional[List[str] | str] = None,
    **dag_kwargs: Any
) -> Callable[[Callable], Any]:
    def decorator(dag_func: Callable) -> Any:
        # 生成原生DAG对象
        dag = original_dag(**dag_kwargs)(dag_func)
        
        if not locks:
            return dag
        
        # 1. 添加获取锁的传感器任务
        lock_task = get_lock.override(task_id="acquire_locks")(locks)
        lock_task.dag = dag
        
        # 2. 添加释放锁的清理任务(teardown类型,确保DAG结束必执行)
        @PythonOperator(task_id="release_locks", teardown=True, dag=dag)
        def unlock_task():
            free_locks(locks)
        
        # 3. 识别DAG中的根任务与叶子任务
        root_tasks = [task for task in dag.tasks if not task.upstream_list]
        leaf_tasks = [task for task in dag.tasks if not task.downstream_list]
        
        # 4. 建立依赖关系
        for root_task in root_tasks:
            lock_task >> root_task
        
        for leaf_task in leaf_tasks:
            leaf_task >> unlock_task
        
        return dag
    return decorator

简化后的DAG使用示例

from datetime import datetime
from airflow.decorators import task

@dag_with_locks(
    locks=["TRANSACTION_VALIDATION"],
    catchup=False,
    tags=["staging"],
    start_date=datetime(2021, 12, 21),
    schedule="*/5 2-17 * * 1-5",
)
def validate_transactions():
    transaction_validation = TransactionValidation()

    @task()
    def get_transactions_to_validate():
        return transaction_validation.get_to_validate_transactions()

    @task()
    def validate_transaction(transaction_id: int):
        return transaction_validation.validate_transactions(transaction_id)

    # 只需编写业务任务,加锁解锁由装饰器自动处理
    t2 = get_transactions_to_validate()
    t3 = validate_transaction.expand(transaction_id=t2)

validate_transactions()

关键优化建议

  • 原子性加锁:当前get_lock逐个检查锁的方式存在竞态风险,建议改用Redis Lua脚本实现原子性的锁检查与获取,避免多个DAG同时抢占部分锁。
  • 锁归属验证:解锁时检查锁的值是否为当前DAG的run_id,防止误释放其他DAG持有的锁。
  • 过期锁清理:添加定时任务清理Redis中过期的锁,避免因DAG异常终止导致锁永久占用。

内容的提问来源于stack exchange,提问作者Lennert De Feyter

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.25 07:20:17