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

Airflow 2.5中PostgresOperator向PythonOperator传数据帧报错

问题分析与修复方案

报错原因

  1. PostgresOperator的XCom返回问题:PostgresOperator执行SELECT语句时,默认返回的是查询结果的元组列表,直接将其传递给PythonOperator的op_args会导致每个任务的结果被拆分成多个位置参数,触发TypeError: too many positional arguments。
  2. 参数名不匹配:get_data函数定义参数为config,但内部使用未定义的sql变量,存在参数名错误。
  3. 动态任务参数传递方式错误: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)

关键修复点说明

  1. 替换PostgresOperator为自定义任务:用@task装饰的get_data函数直接通过PostgresHook执行SQL并返回结构化数据,替代PostgresOperator,更适合获取查询结果并传递给后续任务。
  2. 正确使用动态任务展开:通过task.expand()实现动态任务,get_data.expand(sql=sql_queries)为每个SQL生成独立查询任务,push_data.expand(data_records=data_results)为每个查询结果生成独立处理任务。
  3. 参数与逻辑修正:修正get_data的参数名与内部使用一致,将查询结果转换为字典列表,方便push_data直接遍历处理。
  4. 路径优化:用os.path.join拼接文件路径,避免因系统路径分隔符差异导致的错误。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.07 05:10:43