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

Airflow虚拟环境算子无法加载SQLAlchemy>2.0问题排查

问题

我创建了如下Airflow DAG:

from airflow.operators.empty import EmptyOperator
from airflow.sensors.filesystem import FileSensor
from airflow.decorators import dag, task

def get_filepath(**context):
    base_path = "/home/project"
    current_date = datetime.now() - relativedelta(months=1)
    filename_pattern = f"myfile.zip"
    return str(Path(base_path) / filename_pattern)

@task(task_id='verify_file_exists')
def verify_file_exists(**context):
    path = "/home/project"
    current_date = datetime.now() - relativedelta(months=1)
    filenames = glob(path + f"*.zip")
    if len(filenames) == 0:
        raise Exception("Required file not found.")
    context['task_instance'].xcom_push(key='filename', value=filenames[0])
    return filenames[0]

@task.virtualenv(
    task_id='process_file',
    multiple_outputs=True,
    requirements=["sqlalchemy>2.0.0"],
    system_site_packages=False,
)
def fun_process_file(filename: str):
    from sqlalchemy import create_engine
    from sqlalchemy.engine import URL
    import zipfile
    import pandas as pd
    import time
    from io import BytesIO
    import sqlalchemy

    print(f"{sqlalchemy.__version__=}")


@dag(
    dag_id="mydag",
    start_date=pendulum.datetime(2024, 7, 16),
    schedule="40 7 10 * *",
    default_args=default_args,
    catchup=False,
)
def mydag():
    start = EmptyOperator(task_id='start')

    wait_for_file = FileSensor(
        task_id='wait_for_file',
        filepath=get_filepath(),
        poke_interval=60*60*24,
        mode='reschedule',
        timeout=60*60*24*7,
        soft_fail=False
    )

    end = EmptyOperator(task_id='end')

    file_exists = verify_file_exists()
    process_file = fun_process_file(file_exists)

    start >> wait_for_file >> file_exists >> process_file >> end

dag = mydag()

但始终无法加载SQLAlchemy>2.0版本,运行时打印的sqlalchemy.__version__为'1.4.54'。日志显示:

ERROR: pip's dependency resolver does not currently take into account all the packages that are installed. This behaviour is the source of the following dependency conflicts.

请问这是否意味着无法在Airflow中运行SQLAlchemy>2.0?此外我尝试用Pandas 2.2.3将DataFrame写入数据库时存在兼容性问题,求问原因。


解答

1. 为什么无法安装SQLAlchemy>2.0?

不是Airflow本身无法运行SQLAlchemy>2.0,问题出在你的virtualenv任务依赖配置和pip的依赖解析逻辑:

  • 你的fun_process_file任务中导入了pandas,但requirements参数里没有明确指定pandas版本。pip创建virtualenv环境时,会自动安装满足隐式依赖的pandas版本,部分旧版pandas对SQLAlchemy有<=1.4.x的版本限制,导致pip无法安装你要求的SQLAlchemy>2.0,只能降级到兼容的1.4.54版本。
  • 虽然Airflow 2.6及之前版本依赖SQLAlchemy 1.4.x,但你设置了system_site_packages=False,virtualenv环境是完全隔离的,不会受Airflow主环境依赖影响。

解决方法:
在requirements中同时明确指定兼容的pandas和SQLAlchemy版本,让pip直接安装目标版本,避免依赖解析自动降级:

@task.virtualenv(
    task_id='process_file',
    multiple_outputs=True,
    requirements=["sqlalchemy>2.0.0", "pandas==2.2.3"],
    system_site_packages=False,
)

2. Pandas 2.2.3写入数据库的兼容性问题原因

你遇到的兼容性问题和当前SQLAlchemy版本不匹配直接相关:

  • Pandas 2.2.x开始大量采用SQLAlchemy 2.x的新API(比如Connection.execute()语法、结果集处理逻辑),而SQLAlchemy 1.4.x的部分API已被废弃或行为不同,两者搭配使用会出现语法错误或功能异常。
  • 若后续成功安装SQLAlchemy>2.0,还需注意数据库驱动兼容性:比如PostgreSQL的psycopg2-binary需要>=2.9版本,MySQL的mysql-connector-python需要>=8.0.30版本,才能适配SQLAlchemy 2.x的连接逻辑。

解决方法:

  • 确保virtualenv中安装兼容版本组合:Pandas 2.2.3 + SQLAlchemy>2.0.0 + 适配的数据库驱动。
  • 检查to_sql()调用方式:使用SQLAlchemy引擎实例作为engine参数,避免旧连接字符串格式;若需自定义SQL逻辑,遵循SQLAlchemy 2.x语法规范。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.14 04:12:32