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默认采用逐行插入模式,未开启批量提交,不仅写入速度慢,还容易触发连接超时、事务锁等问题 - 连接管理不规范:未使用上下文管理器管理数据库连接,任务异常抛出时连接不会自动释放,长期运行会占满数据库连接池
修复方案
- 复用Airflow内置连接配置,不要硬编码数据库地址、账号密码,从Airflow Connections中读取
bd_proj_postgres的配置生成SQLAlchemy连接串,保证和PostgresOperator的网络配置完全一致 - 修正字段类型,将
date类型的日期值转换为datetime类型,匹配表结构的timestamp字段定义 - 开启批量写入参数,设置合理的chunksize,使用multi批量提交模式提升写入性能
- 用上下文管理器管理数据库连接,任务结束/异常时自动释放连接
修复后写入函数代码
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
相关产品推荐
相关产品推荐

