能否配置Airflow Sensor(如@task.sensor)监听GitHub PR合并触发DAG运行?
用Airflow @task.sensor监听GitHub PR合并事件(无需GitHub Actions)
完全可以实现,核心思路是通过GitHub API定期轮询目标仓库的PR状态,用Airflow的@task.sensor装饰器封装这个轮询逻辑,全程不需要依赖GitHub Actions。具体实现步骤如下:
1. 准备GitHub认证
生成一个GitHub个人访问令牌(PAT),权限勾选repo(确保能访问目标仓库的PR数据)。然后在Airflow控制台的连接中新增一个HTTP类型的连接:
- 连接ID设为
github_api_conn - 主机填
api.github.com - 密码字段填入生成的PAT
同时在Airflow变量中添加三个变量:
github_repo_owner:仓库所属用户名/组织名github_repo_name:目标仓库名称last_merged_pr_id:初始值设为0,用来记录上次触发的合并PR ID,避免重复触发
2. 编写传感器DAG
用@task.sensor装饰器实现轮询逻辑,设置轮询间隔、超时时间,在poke函数中调用GitHub API验证PR合并状态:
from airflow.decorators import dag, task from airflow.models import Variable from airflow.providers.http.hooks.http import HttpHook from datetime import datetime, timedelta default_args = { 'owner': 'airflow', 'retries': 1, 'retry_delay': timedelta(minutes=5), } @dag( default_args=default_args, schedule=None, # 传感器触发,无需定时调度 start_date=datetime(2024, 1, 1), catchup=False, ) def github_pr_merge_monitor(): @task.sensor(poke_interval=300, timeout=86400, mode='poke') def check_pr_merge_status(): # 初始化GitHub API钩子 http_hook = HttpHook(http_conn_id='github_api_conn', method='GET') repo_owner = Variable.get('github_repo_owner') repo_name = Variable.get('github_repo_name') last_pr_id = int(Variable.get('last_merged_pr_id', default_var=0)) # 调用GitHub API获取已关闭的PR列表 api_endpoint = f'/repos/{repo_owner}/{repo_name}/pulls?state=closed' response = http_hook.run(api_endpoint) pr_list = response.json() # 过滤出已合并且ID大于上次记录的PR new_merged_prs = [pr for pr in pr_list if pr['merged'] and pr['id'] > last_pr_id] if new_merged_prs: # 更新最后合并的PR ID,避免重复触发 latest_pr_id = max(pr['id'] for pr in new_merged_prs) Variable.set('last_merged_pr_id', str(latest_pr_id)) return True return False # 传感器触发后续业务任务 pr_merge_trigger = check_pr_merge_status() @task def execute_post_merge_workflow(): print("检测到新的PR合并事件,执行后续工作流") # 这里可以添加具体的业务逻辑,比如代码部署、数据同步等 pr_merge_trigger >> execute_post_merge_workflow() github_pr_merge_monitor()
3. 关键注意事项
- API速率限制:GitHub API免费额度为每小时5000次请求,建议将
poke_interval设为5-10分钟(300-600秒),避免触发限制 - 状态持久化:用Airflow Variable存储最后合并PR ID时,若Airflow是多实例部署,需注意变量的并发读写问题,也可以改用数据库表来存储状态
- 传感器模式:默认用
poke模式,若需要长时间监听且节省资源,可改为reschedule模式(传感器会在等待时释放worker资源) - 权限验证:确保PAT拥有目标仓库的
repo权限,否则会返回API访问错误
内容的提问来源于stack exchange,提问作者odonnry
相关产品推荐
相关产品推荐

