如何为Airflow PythonOperator传递命令行参数并解决调用报错
我有一个可接收命令行参数的Python脚本,需要在Airflow DAG中通过PythonOperator调用该脚本的test方法。使用op_args传递参数时出现如下报错:
{standard_task_runner.py:107} ERROR - Failed to execute job 56
for task test (test() takes 1 positional argument but 4 were given; 544)
直接在命令行执行python test.py -f Hello -b World,或是在其他函数中调用test(['-f','Hello','-b','World'])都能正常运行,但在Airflow DAG中使用op_args = ['-f','Hello','-b','World']传递参数时却失败了。
我的DAG代码如下:
from airflow.operators import PythonOperator from airflow.models import DAG from datetime import datetime, timedelta from test import test default_args = { 'owner': 'none', 'retries': 5, 'retry_delay': timedelta(minutes=5) } with DAG( dag_id='test', default_args = default_args, description =' Still in process', start_date=datetime(2022,9,13), schedule_interval='@weekly' ) as DAG: task = PythonOperator( task_id='test', python_callable=test, op_args = ['-f','Hello','-b','World'] ) task
test函数代码如下:
import argparse import sys def test(args): parser = argparse.ArgumentParser(description='Description of your program') parser.add_argument('-f','--foo', help='Description for foo argument', required=True) parser.add_argument('-b','--bar', help='Description for bar argument', required=True) args = parser.parse_args(args) if args.foo == 'Hello': # code here print('hello') if args.bar == 'World': # code here print('world') if __name__ =='__main__': test(sys.argv[1:])
问题根源在于Airflow的op_args传递逻辑:它会把列表中的每个元素作为单独的位置参数传给指定的python_callable。而你的test函数仅接受一个参数(即包含所有命令行参数的列表),所以当你传入4个元素的列表时,相当于给test传了4个位置参数,自然触发参数不匹配的报错。
有两种简单的解决方式:
方式一:嵌套列表传递参数
把op_args改成嵌套列表,让整个参数列表作为单个元素传递给test函数:
task = PythonOperator( task_id='test', python_callable=test, op_args = [['-f','Hello','-b','World']] # 外层列表包裹内层参数列表 )
这样test函数会收到内层的['-f','Hello','-b','World'],和你直接调用test(sys.argv[1:])的形式完全一致,argparse可以正常解析。
方式二:修改test函数接受可变参数
调整test函数的定义,让它接受可变数量的位置参数,再打包后传给argparse:
import argparse import sys def test(*args): # 用*args接收所有传入的位置参数 parser = argparse.ArgumentParser(description='Description of your program') parser.add_argument('-f','--foo', help='Description for foo argument', required=True) parser.add_argument('-b','--bar', help='Description for bar argument', required=True) args = parser.parse_args(args) # args是元组,parse_args可以直接处理 if args.foo == 'Hello': print('hello') if args.bar == 'World': print('world') if __name__ =='__main__': test(*sys.argv[1:]) # 命令行调用时也用*解包传递
修改后,原有的op_args = ['-f','Hello','-b','World']可以正常工作,因为*args会把传入的4个参数打包成一个元组,而argparse.ArgumentParser.parse_args()支持处理元组类型的输入。
内容的提问来源于stack exchange,提问作者Time

