Apache Airflow DAG导入报错FileNotFoundError:文件未创建却被提前读取
Airflow DAG导入报错问题
我在DAG中配置了两个任务:
- 任务1:通过SimpleHttpOperator从API拉取数据,在Worker节点创建
/tmp/data目录,并将API返回数据写入/tmp/data/test_data.json文件,该任务执行正常,可在Worker节点访问到目标文件。 - 添加任务2后,DAG无法导入,提示报错
FileNotFoundError: [Errno 2] No such file or directory,原因是/tmp/data/test_data.json文件不存在。但该文件需待任务1执行后才会创建,疑惑为何任务运行前就触发报错;单独编写读写文件的Python代码可正常运行,推测是对Airflow特性不了解导致。
任务1代码
get_data = SimpleHttpOperator( task_id='get_data', method='GET', endpoint='endpoint', http_conn_id='api', headers={'Authorization': 'Bearer'}, response_check=lambda response: _handle_response(response), dag=dag )
_handle_response函数
def _handle_response(response): print(response.status_code) pathlib.Path("/tmp/data").mkdir(parents=True, exist_ok=True) with open("/tmp/data/test_data.json","wb") as f: f.write(response.content) return True
任务2代码
read_data = PythonOperator( task_id='read_data', python_callable=_data_to_read("/tmp/data/test_data.json"), dag=dag )
_data_to_read函数
def _data_to_read(xcom): with open(xcom) as json_file: data = json.load(json_file)
问题根因
Task 2的python_callable参数直接调用了_data_to_read函数,而非传入函数引用,导致Airflow解析DAG时就执行了文件读取操作,此时任务1尚未运行,文件未创建,从而触发导入错误。
内容的提问来源于stack exchange,提问作者Sprant
相关产品推荐
相关产品推荐

