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
相关产品推荐
相关产品推荐

