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

如何在Airflow 2.6.1中使用RdsHook查询RDS数据?

使用Airflow 2.6.1的RdsHook查询RDS数据操作流程

前置确认

确保你已安装适配Airflow 2.6.1的AWS Providers包:

pip install apache-airflow-providers-amazon>=8.0.0

完整操作步骤

  • 导入所需模块
    根据你的RDS数据库引擎(比如PostgreSQL用psycopg2,MySQL用pymysql),导入对应驱动和RdsHook:

    from airflow.providers.amazon.aws.hooks.rds import RdsHook
    # PostgreSQL示例,MySQL替换为import pymysql
    import psycopg2
    from psycopg2.extras import DictCursor
    
  • 初始化RdsHook并获取连接参数
    通过你创建的rds_conn连接ID实例化钩子,提取RDS的连接配置:

    # 初始化RdsHook
    rds_hook = RdsHook(rds_conn_id="rds_conn")
    # 提取Airflow连接中的配置信息
    conn_config = rds_hook.get_connection("rds_conn")
    
  • 建立数据库连接并执行查询
    用提取到的参数构建数据库连接,执行SQL查询并处理结果:

    # PostgreSQL示例,MySQL替换为pymysql.connect
    with psycopg2.connect(
        host=conn_config.host,
        port=conn_config.port,
        dbname=conn_config.schema,
        user=conn_config.login,
        password=conn_config.password
    ) as conn:
        with conn.cursor(cursor_factory=DictCursor) as cursor:
            # 执行查询SQL,替换为你的实际语句
            cursor.execute("SELECT * FROM your_target_table LIMIT 10;")
            # 获取查询结果
            results = cursor.fetchall()
            # 按需处理结果,比如打印或写入存储
            for row in results:
                print(row)
    

注意事项

  • 若使用MySQL RDS,需将psycopg2相关代码替换为pymysql,连接参数中的dbname改为database
  • 建议将上述逻辑封装到Airflow的PythonOperator中,作为DAG的任务节点执行
  • 确保Airflow Worker节点的网络能访问RDS实例(检查安全组、网络ACL配置)

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.12 05:38:15