Airflow 2.5中PostgresOperator向PythonOperator传数据帧报错
问题分析与修复方案
报错原因
- PostgresOperator的XCom返回问题:PostgresOperator执行SELECT语句时,默认返回的是查询结果的元组列表,直接将其传递给PythonOperator的
op_args会导致每个任务的结果被拆分成多个位置参数,触发TypeError: too many positional arguments。 - 参数名不匹配:
get_data函数定义参数为config,但内部使用未定义的sql变量,存在参数名错误。 - 动态任务参数传递方式错误:PythonOperator的
expand(op_args=...)要求每个元素是传递给python_callable的参数元组,直接传入XCom返回的结果集会导致参数拆分异常。
修复后的完整代码
import os import json import logging from datetime import datetime, timedelta from airflow.decorators import task from airflow import DAG from airflow.hooks.postgres_hook import PostgresHook from airflow.models import Variable BASE_DIR = Variable.get("BASE_DIR") log = logging.getLogger("airflow") DEFAULT_ARGS = { 'owner': 'airflow', 'depends_on_past': False, 'start_date': datetime(2022, 12, 16), 'email': [''], 'email_on_failure': False, 'email_on_retry': False, 'retries': 0, 'retry_delay': timedelta(minutes=1) } def parse_columns(arr): str1 = "," if isinstance(arr, dict): return str1.join(arr.keys()) @task def generate_sql_queries(): config_filepath = f"{BASE_DIR}/sd/retl/map/" queries = [] for filename in os.listdir(config_filepath): with open(os.path.join(config_filepath, filename)) as f: config = json.load(f) cols = parse_columns(config['properties_map']) query = f"SELECT {cols} FROM {config['source_stream']['schema']}.{config['source_stream']['name']}" queries.append(query) return queries @task def get_data(sql): """执行SQL查询并返回结构化数据""" pg_hook = PostgresHook(postgres_conn_id='DWH_PROD') df = pg_hook.get_pandas_df(sql=sql) log.info(f"查询到{len(df)}条数据") return df.to_dict('records') # 转换为字典列表,方便后续遍历处理 @task def push_data(data_records): """遍历结果集执行后续处理逻辑""" for record in data_records: log.info(f"处理数据: {record}") # 在这里添加你的后续处理逻辑,比如推送数据到外部系统、数据清洗等 with DAG( dag_id='SD_rETL', default_args=DEFAULT_ARGS, schedule_interval='@daily', start_date=datetime(year=2022, month=2, day=1), catchup=False ) as dag: sql_queries = generate_sql_queries() # 为每个SQL生成独立的查询任务 data_results = get_data.expand(sql=sql_queries) # 为每个查询结果生成独立的处理任务 push_data.expand(data_records=data_results)
关键修复点说明
- 替换PostgresOperator为自定义任务:用
@task装饰的get_data函数直接通过PostgresHook执行SQL并返回结构化数据,替代PostgresOperator,更适合获取查询结果并传递给后续任务。 - 正确使用动态任务展开:通过
task.expand()实现动态任务,get_data.expand(sql=sql_queries)为每个SQL生成独立查询任务,push_data.expand(data_records=data_results)为每个查询结果生成独立处理任务。 - 参数与逻辑修正:修正
get_data的参数名与内部使用一致,将查询结果转换为字典列表,方便push_data直接遍历处理。 - 路径优化:用
os.path.join拼接文件路径,避免因系统路径分隔符差异导致的错误。
内容的提问来源于stack exchange,提问作者Andriychuk Kirill
相关产品推荐
相关产品推荐

