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

Airflow从MySQL迁移数据到CSV文件时报函数对象不可迭代错误

Airflow MySQL导出CSV报function object is not iterable的解决方法

错误根因

触发该报错的直接原因是:MySqlOperator的sql参数要求传入SQL字符串/包含SQL字符串的可迭代对象,但你直接传入了extract_data函数对象,运算符尝试遍历该参数执行SQL时就会抛出函数不可迭代的错误。

除此之外你的代码还存在多个隐藏问题:

  • extract_data函数中的mysql_conn_id没有定义,直接传给pd.read_sql会触发变量未定义报错
  • 不同Operator的任务是独立运行的进程,extract任务中生成的df变量无法直接在dbpost_process任务中调用,会触发变量未定义报错
  • MySqlOperator的定位是执行SQL命令(如建表、插入、更新),不会直接返回查询结果的DataFrame对象,属于运算符使用错误。

完整解决方案

方案1:单PythonOperator完成全流程(推荐,小数据量场景无跨任务传值问题)

from airflow import DAG
from datetime import datetime
from airflow.operators.python_operator import PythonOperator
import pandas as pd
from airflow.hooks.mysql_hook import MySqlHook

default_args = {"owner":"airflow","start_date":datetime(2021,7,10)}

def export_mysql_to_csv():
    # 调用Airflow内置MySQLHook获取数据库连接,无需硬编码账号密码
    mysql_hook = MySqlHook(mysql_conn_id="mysql_db1")
    conn = mysql_hook.get_conn()
    # 查询转DataFrame后直接导出CSV
    df = pd.read_sql('SELECT * FROM new_table', conn)
    df.to_csv('~/op_files/sample3.csv', index=False)

with DAG(dag_id="dbworkflow2",default_args=default_args,schedule_interval=None) as dag:
    export_task = PythonOperator(
        task_id = "export_mysql_to_csv",
        python_callable = export_mysql_to_csv
    )

方案2:拆分提取和导出两个任务(适合需要中间逻辑处理的场景)

from airflow import DAG
from datetime import datetime
from airflow.operators.python_operator import PythonOperator
import pandas as pd
from airflow.hooks.mysql_hook import MySqlHook

default_args = {"owner":"airflow","start_date":datetime(2021,7,10)}

def extract_data(**context):
    mysql_hook = MySqlHook(mysql_conn_id="mysql_db1")
    conn = mysql_hook.get_conn()
    df = pd.read_sql('SELECT * FROM new_table', conn)
    # 序列化DataFrame后推送到XCom实现跨任务传值
    context['ti'].xcom_push(key='mysql_query_result', value=df.to_json(orient='records'))

def dbpost_process(**context):
    # 从XCom拉取查询结果转回DataFrame
    df_json = context['ti'].xcom_pull(key='mysql_query_result', task_ids='extract')
    df = pd.read_json(df_json, orient='records')
    df.to_csv('~/op_files/sample3.csv', index=False)

with DAG(dag_id="dbworkflow2",default_args=default_args,schedule_interval=None) as dag:
    extract = PythonOperator(
        task_id='extract',
        python_callable = extract_data,
        provide_context=True
    )
    dbpost_process = PythonOperator(
        task_id = "dbpost_process",
        python_callable = dbpost_process,
        provide_context=True
    )   
    extract >> dbpost_process

内容的提问来源于stack exchange,提问作者Anil Krishna Parvathaneni

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.01 14:45:02