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

能否配置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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.03 14:15:35