如何编写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
相关产品推荐
相关产品推荐

