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

如何通过Airflow DAG跨Oracle数据库连接执行表Join查询

跨Oracle库表Join的Airflow实现方案

方案一:使用Oracle数据库链接(DBLINK)

这是效率最优的方案,Join操作直接在数据库层面执行。

  • 步骤1:在其中一个Oracle库(如db1,对应Airflow连接ID c1)创建指向db2的数据库链接
    -- 登录db1执行该SQL
    CREATE DATABASE LINK db2_link
    CONNECT TO db2_username IDENTIFIED BY db2_password
    USING '(DESCRIPTION =
      (ADDRESS_LIST =
        (ADDRESS = (PROTOCOL = TCP)(HOST = db2_host)(PORT = db2_port))
      )
      (CONNECT_DATA =
        (SID = db2_sid)
      )
    )';
    
  • 步骤2:编写跨库Join的SQL文件(示例命名为join_query.sql)
    SELECT t1.*, t2.*
    FROM table1 t1
    JOIN table2@db2_link t2
      ON t1.join_key = t2.join_key;
    
  • 步骤3:修改Airflow函数并编写DAG
    只需用db1的连接执行SQL即可,无需同时调用两个连接:
    from airflow import DAG
    from airflow.providers.oracle.hooks.oracle import OracleHook
    from airflow.operators.python import PythonOperator
    from datetime import datetime
    import logging
    
    logger = logging.getLogger(__name__)
    
    ORACLE_CONN_ID1 = 'c1'  # 对应db1的Airflow连接ID
    
    def extract_from_oracle(extract_sql):
        logger.info(f'执行SQL文件: {extract_sql}')
        try:
            oracle_hook = OracleHook(ORACLE_CONN_ID1)
            with oracle_hook.get_conn() as conn:
                with conn.cursor() as cursor:
                    with open(extract_sql, 'r') as f:
                        sql = f.read()
                    cursor.execute(sql)
                    # 可选:获取查询结果
                    results = cursor.fetchall()
                    logger.info(f'查询返回 {len(results)} 条数据')
            return results
        except Exception as e:
            logger.error(f'查询失败: {str(e)}')
            raise
    
    with DAG(
        'oracle_cross_db_join',
        start_date=datetime(2024, 1, 1),
        schedule_interval=None,
        catchup=False
    ) as dag:
        join_task = PythonOperator(
            task_id='cross_db_join_query',
            python_callable=extract_from_oracle,
            op_kwargs={'extract_sql': '/path/to/join_query.sql'}
        )
    

方案二:Airflow拉取数据后在应用层Join

适合无法创建DBLINK的场景,用Pandas在Python Operator中实现Join:

  • 步骤1:编写拉取数据并Join的函数及DAG
    from airflow import DAG
    from airflow.providers.oracle.hooks.oracle import OracleHook
    from airflow.operators.python import PythonOperator
    from datetime import datetime
    import logging
    import pandas as pd
    
    logger = logging.getLogger(__name__)
    
    ORACLE_CONN_ID1 = 'c1'  # db2的Airflow连接ID
    ORACLE_CONN_ID2 = 'c2'  # db2的Airflow连接ID
    
    def cross_db_join():
        try:
            # 从db1拉取table1数据
            hook1 = OracleHook(ORACLE_CONN_ID1)
            df1 = hook1.get_pandas_df("SELECT * FROM table1")
            logger.info(f'从db1获取 {len(df1)} 条数据')
    
            # 从db2拉取table2数据
            hook2 = OracleHook(ORACLE_CONN_ID2)
            df2 = hook2.get_pandas_df("SELECT * FROM table2")
            logger.info(f'从db2获取 {len(df2)} 条数据')
    
            # 执行Inner Join,可根据需求修改Join类型
            joined_df = pd.merge(df1, df2, on='join_key', how='inner')
            logger.info(f'Join后得到 {len(joined_df)} 条数据')
    
            # 可选:将结果写入目标库或存储
            # hook1.insert_rows(table='target_table', rows=joined_df.values.tolist())
            return joined_df.to_json()
        except Exception as e:
            logger.error(f'Join失败: {str(e)}')
            raise
    
    with DAG(
        'oracle_cross_db_join_pandas',
        start_date=datetime(2024, 1, 1),
        schedule_interval=None,
        catchup=False
    ) as dag:
        join_task = PythonOperator(
            task_id='pandas_cross_db_join',
            python_callable=cross_db_join
        )
    

原代码问题说明

你提供的代码存在逻辑错误:

  • oracle_conn_id 是字符串列表,con=0 是整数,while con in oracle_conn_id 永远不成立,循环不会执行
  • 跨库Join无法通过两个独立数据库连接实现,数据库层面无法识别另一个连接的表,必须通过DBLINK或应用层聚合数据

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.23 06:05:05