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

如何用Airflow AwsGlueCatalogPartitionSensor监测Glue未知分区更新?

解决Glue Catalog分区更新触发Airflow DAG的方案

核心思路

不用死磕AwsGlueCatalogPartitionSensor必须指定分区表达式的限制,换个思路——定期扫描目标表的分区元数据,对比分区的最后修改时间,不管上游更新的是新分区还是一个月前的旧分区,只要分区的修改时间发生变化,就触发后续DAG任务。

具体实现步骤

1. 自定义分区更新传感器

直接写个自定义传感器(或者用PythonOperator实现检查逻辑),核心逻辑是对比分区的修改时间:

  • 调用Glue API拉取目标表最近N天的所有分区(N设为30,覆盖上游可能修改的旧分区范围)
  • 读取Airflow变量/外部数据库中存储的「上次扫描的分区修改时间记录」
  • 对比当前分区的修改时间,找出新增或已更新的分区
  • 存在更新则标记传感器成功,触发后续任务;否则继续等待

2. 代码示例

from airflow.providers.amazon.aws.hooks.glue import GlueCatalogHook
from airflow.sensors.base import BaseSensorOperator
from airflow.models import Variable
import datetime

class GluePartitionUpdateSensor(BaseSensorOperator):
    def __init__(self, table_name, database_name, lookback_days=30, aws_conn_id="aws_default", **kwargs):
        super().__init__(**kwargs)
        self.table_name = table_name
        self.database_name = database_name
        self.lookback_days = lookback_days
        self.glue_hook = GlueCatalogHook(aws_conn_id=aws_conn_id)

    def poke(self, context):
        # 计算扫描的起始日期,覆盖上游可能修改的旧分区
        cutoff_date = (datetime.datetime.now() - datetime.timedelta(days=self.lookback_days)).strftime("%Y-%m-%d")
        # 获取指定日期范围内的所有分区
        partitions = self.glue_hook.get_partitions(
            database_name=self.database_name,
            table_name=self.table_name,
            expression=f"business_date >= '{cutoff_date}'"
        )
        # 整理当前分区的修改时间:用business_date+source作为唯一键
        current_partition_mods = {}
        for part in partitions:
            part_key = f"{part['Values'][0]}_{part['Values'][1]}"
            # 优先取存储参数里的last_modified,没有就用创建时间
            mod_time = part['StorageDescriptor']['Parameters'].get('last_modified', part['CreationTime'].isoformat())
            current_partition_mods[part_key] = mod_time
        
        # 读取上次保存的分区修改记录
        last_recorded = Variable.get("glue_partition_mod_records", default_var={}, deserialize_json=True)
        
        # 检查是否有更新
        has_update = False
        for key, mod_time in current_partition_mods.items():
            if key not in last_recorded or mod_time != last_recorded[key]:
                has_update = True
                break
        
        # 有更新则保存最新记录,返回True触发后续任务
        if has_update:
            Variable.set("glue_partition_mod_records", current_partition_mods, serialize_json=True)
            return True
        return False

3. DAG配置要点

  • 把这个自定义传感器作为DAG的起始任务
  • 设置DAG的调度间隔为短周期(比如5分钟),保证能及时捕获上游的分区更新
  • 传感器的poke_interval可以设为300秒(和调度间隔匹配),避免频繁调用Glue API

4. 实用注意事项

  • 确保Airflow使用的AWS连接拥有glue:GetPartitions权限,否则无法读取分区元数据
  • 如果分区数量过多,Airflow变量存储有长度限制,建议改用Airflow元数据库或外部轻量数据库(比如SQLite)存储分区修改记录
  • 根据上游更新频率调整lookback_days和调度间隔:比如上游一天更新5次,调度间隔设5分钟;若上游很少修改超过30天的分区,也可以缩小lookback_days减少API调用量

内容的提问来源于stack exchange,提问作者Prateek Pathak

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.19 16:05:01