如何通过Airflow DAG跨Oracle数据库连接执行表Join查询
跨Oracle库表Join的Airflow实现方案
方案一:使用Oracle数据库链接(DBLINK)
这是效率最优的方案,Join操作直接在数据库层面执行。
- 步骤1:在其中一个Oracle库(如
db1,对应Airflow连接IDc1)创建指向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
相关产品推荐
相关产品推荐

