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

Airflow PythonOperator通过SQLAlchemy写入PostgreSQL数据失败排查

Airflow PostgreSQL数据同步DAG故障修复

问题描述

  • 开发Benzinga盘前股票行情爬取同步DAG时,PostgresOperator可正常完成目标表创建操作,但负责写入DataFrame的PythonOperator任务执行失败
  • 已尝试重启PostgreSQL服务、重启运行设备,参考公开技术社区方案均未解决问题
  • 核心需求:通过SQLAlchemy连接PostgreSQL,将爬取得到的行情DataFrame一次性批量写入目标表
  • 原始复现代码如下:
from airflow import DAG
from airflow.providers.postgres.operators.postgres import PostgresOperator
from airflow.operators.python_operator import PythonOperator
import pandas as pd
from sqlalchemy import create_engine
from datetime import timedelta, datetime
from airflow.utils.dates import days_ago

default_args = {
    'owner': 'rudresh_mehta',
    'retries': 0,
    'retry_delay': timedelta(minutes=1)
}

def fetch_premarket_table(ti):
    url = "https://www.benzinga.com/premarket/"
    tables = pd.read_html(url)
    premarket_df = tables[5].to_json()
    ti.xcom_push(key="premarket_df",value=premarket_df)

def store_premarket_table(**context):
    import psycopg2
    ti = context["ti"]
    json_str = ti.xcom_pull(key="premarket_df", task_ids=["fetch_bezinga_premarket_table_id"])

    print("json",json_str)
    df = pd.read_json(json_str[0])
    print("df================",df)

    current_time = datetime.today()
    today = current_time.date()
    time_today = datetime.now().time()
    yesterday = today - timedelta(days=1)
    print("after today===========================")
    # 硬编码连接配置
    conn_string = 'postgresql://postgres:root@localhost:5432/bd_proj'
    print("============",conn_string)
    db = create_engine(conn_string, pool_size=10, max_overflow=20)
    connection = db.connect()
    connection.autocommit = True
    df["date_time"] = yesterday
    # 逐行写入
    df.to_sql(name='bezinga_premarket_table', con=connection, if_exists='append', index=False)
    connection.close()
    print("Success")
    return 200


with DAG(
    dag_id='dag_with_modulev01',
    default_args=default_args,
    start_date=datetime(2022, 7, 8),
    schedule_interval='8 5 * * *') as dag:

    create_premarket_tb = PostgresOperator(
        task_id='create_premarket_table_id',
        postgres_conn_id='bd_proj_postgres',
        sql="""
            create table if not exists bezinga_premarket_table(
                stock_ticker character(8),
                stock_name character varying,
                stock_price real,
                change_perc character(10),
                volume character varying,
                date_time timestamp,
                primary key(stock_ticker, date_time));
        """)
    fetching_premarket_tb = PythonOperator(task_id="fetch_bezinga_premarket_table_id",
        python_callable=fetch_premarket_table,dag=dag )

    store_premarket_tb = PythonOperator(
        task_id='store_premarket_table_id', provide_context=True,
        python_callable=store_premarket_table, dag=dag)

    create_premarket_tb >> fetching_premarket_tb >> store_premarket_tb

报错截图

核心故障点

  • 连接地址配置错误:PostgresOperator使用Airflow内置的bd_proj_postgres连接配置可正常访问数据库,但Python函数中硬编码的localhost/127.0.0.1指向的是Airflow Worker自身的运行环境(如果是Docker部署、分布式部署场景,Worker和PostgreSQL不在同一网络/主机下,该地址根本无法连通数据库)
  • 字段类型不匹配:表结构中date_time为timestamp类型,但代码中赋值的yesterday是date类型,类型不兼容会直接触发写入报错
  • 写入逻辑低效:pandas.to_sql默认采用逐行插入模式,未开启批量提交,不仅写入速度慢,还容易触发连接超时、事务锁等问题
  • 连接管理不规范:未使用上下文管理器管理数据库连接,任务异常抛出时连接不会自动释放,长期运行会占满数据库连接池

修复方案

  1. 复用Airflow内置连接配置,不要硬编码数据库地址、账号密码,从Airflow Connections中读取bd_proj_postgres的配置生成SQLAlchemy连接串,保证和PostgresOperator的网络配置完全一致
  2. 修正字段类型,将date类型的日期值转换为datetime类型,匹配表结构的timestamp字段定义
  3. 开启批量写入参数,设置合理的chunksize,使用multi批量提交模式提升写入性能
  4. 用上下文管理器管理数据库连接,任务结束/异常时自动释放连接

修复后写入函数代码

from airflow.hooks.base import BaseHook

def store_premarket_table(**context):
    import pandas as pd
    from sqlalchemy import create_engine
    from datetime import timedelta, datetime

    ti = context["ti"]
    json_str = ti.xcom_pull(key="premarket_df", task_ids=["fetch_bezinga_premarket_table_id"])
    df = pd.read_json(json_str[0])

    # 从Airflow连接配置获取数据库信息,避免硬编码
    conn = BaseHook.get_connection("bd_proj_postgres")
    conn_string = f"postgresql://{conn.login}:{conn.password}@{conn.host}:{conn.port}/{conn.schema}"
    db = create_engine(conn_string, pool_size=10, max_overflow=20)

    # 修正字段类型,匹配timestamp格式
    yesterday = datetime.today() - timedelta(days=1)
    df["date_time"] = yesterday

    # 批量写入,自动管理连接
    with db.connect() as connection:
        df.to_sql(
            name='bezinga_premarket_table',
            con=connection,
            if_exists='append',
            index=False,
            chunksize=1000,
            method='multi'
        )
        connection.commit()
    print("Success")
    return 200

额外排查验证点

  • 网络连通性验证:进入Airflow Worker的实际运行环境(容器/物理机),直接使用psql客户端测试数据库连通性,不要在宿主机环境测试
  • 字段长度校验:爬取得到的stock_ticker字段长度不要超过表定义的character(8)限制,change_perc字段长度不要超过character(10)限制,长度溢出也会触发写入报错
  • 依赖版本校验:确认SQLAlchemy、psycopg2-binary、pandas三个包的版本兼容,不要使用过低或过高的跨大版本依赖

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.27 01:45:42