Airflow DAG执行报错:'S3Hook'对象无'download_file'属性
解决Airflow S3Hook报错:'S3Hook' object has no attribute 'download_file'
问题详情
运行Airflow DAG的download_from_s3任务时触发报错,核心错误为:'S3Hook' object has no attribute 'download_file',完整报错日志如下:
[2023-02-05 11:32:43,293] {{taskinstance.py:887}} INFO - Executing <Task(PythonOperator): download_from_s3> on 2023-02-05T11:32:34.016335+00:00 [2023-02-05 11:32:43,299] {{standard_task_runner.py:53}} INFO - Started process 87503 to run task [2023-02-05 11:32:43,474] {{logging_mixin.py:112}} INFO - Running %s on host %s <TaskInstance: s3_download.download_from_s3 2023-02-05T11:32:34.016335+00:00 [running]> 67c7842be21b [2023-02-05 11:32:43,555] {{taskinstance.py:1128}} ERROR - 'S3Hook' object has no attribute 'download_file' Traceback (most recent call last): File "/usr/local/lib/python3.7/site-packages/airflow/models/taskinstance.py", line 966, in _run_raw_task result = task_copy.execute(context=context) File "/usr/local/lib/python3.7/site-packages/airflow/operators/python_operator.py", line 113, in execute return_value = self.execute_callable() File "/usr/local/lib/python3.7/site-packages/airflow/operators/python_operator.py", line 118, in execute_callable return self.python_callable(*self.op_args, **self.op_kwargs) File "/usr/local/airflow/dags/dwnld_frm_awss3.py", line 12, in download_from_s3 file_name = hook.download_file(key=key, bucket_name=bucket_name, local_path=local_path) AttributeError: 'S3Hook' object has no attribute 'download_file' [2023-02-05 11:32:43,570] {{taskinstance.py:1185}} INFO - Marking task as FAILED.dag_id=s3_download, task_id=download_from_s3, execution_date=20230205T113234, start_date=20230205T113243, end_date=20230205T113243
用户编写的DAG代码:
import os from datetime import datetime from airflow.models import DAG from airflow.operators.python_operator import PythonOperator from airflow.hooks.S3_hook import S3Hook from airflow.contrib.hooks.aws_hook import AwsHook # Function of the DAG def download_from_s3(key: str, bucket_name: str, local_path: str) -> str: hook = S3Hook('my_conn_S3') file_name = hook.download_file(key=key, bucket_name=bucket_name, local_path=local_path) return file_name with DAG( dag_id='s3_download', schedule_interval='@daily', start_date=datetime(2023, 2, 4), catchup=False ) as dag: task_download_from_s3 = PythonOperator( task_id='download_from_s3', python_callable=download_from_s3, op_kwargs={ 'key': 'sample.txt', 'bucket_name': 'airflow-sample-s3-bucket', 'local_path': '/usr/local/airflow/' } )
错误原因
报错核心是Airflow版本差异:
- 当前使用的是Airflow 1.x版本,该版本的
S3Hook(导入路径from airflow.hooks.S3_hook import S3Hook)并没有download_file方法; download_file是Airflow 2.x版本新增的方法,且2.x版本的S3Hook导入路径已变更为from airflow.providers.amazon.aws.hooks.s3 import S3Hook。
解决方法
方法1:适配Airflow 1.x版本(无需升级)
修改download_from_s3函数,通过get_key获取S3对象后执行下载:
import os from datetime import datetime from airflow.models import DAG from airflow.operators.python_operator import PythonOperator from airflow.hooks.S3_hook import S3Hook from airflow.contrib.hooks.aws_hook import AwsHook def download_from_s3(key: str, bucket_name: str, local_path: str) -> str: hook = S3Hook('my_conn_S3') # 获取S3中的文件对象 s3_key = hook.get_key(key, bucket_name=bucket_name) # 拼接本地完整文件路径 local_file_path = os.path.join(local_path, key) # 执行下载操作 s3_key.download_file(local_file_path) return local_file_path with DAG( dag_id='s3_download', schedule_interval='@daily', start_date=datetime(2023, 2, 4), catchup=False ) as dag: task_download_from_s3 = PythonOperator( task_id='download_from_s3', python_callable=download_from_s3, op_kwargs={ 'key': 'sample.txt', 'bucket_name': 'airflow-sample-s3-bucket', 'local_path': '/usr/local/airflow/' } )
方法2:升级到Airflow 2.x版本
- 升级Airflow并安装AWS provider包:
pip install apache-airflow==2.x.x apache-airflow-providers-amazon
(将2.x.x替换为具体的2.x版本号,比如2.5.3)
- 修改DAG中的
S3Hook导入路径,保留download_file方法的使用:
import os from datetime import datetime from airflow.models import DAG from airflow.operators.python_operator import PythonOperator # 替换为2.x版本的S3Hook导入路径 from airflow.providers.amazon.aws.hooks.s3 import S3Hook from airflow.providers.amazon.aws.hooks.base_aws import AwsHook def download_from_s3(key: str, bucket_name: str, local_path: str) -> str: hook = S3Hook('my_conn_S3') # 2.x版本中download_file方法可用 file_name = hook.download_file(key=key, bucket_name=bucket_name, local_path=local_path) return file_name with DAG( dag_id='s3_download', schedule_interval='@daily', start_date=datetime(2023, 2, 4), catchup=False ) as dag: task_download_from_s3 = PythonOperator( task_id='download_from_s3', python_callable=download_from_s3, op_kwargs={ 'key': 'sample.txt', 'bucket_name': 'airflow-sample-s3-bucket', 'local_path': '/usr/local/airflow/' } )
内容的提问来源于stack exchange,提问作者Austin Jackson
相关产品推荐
相关产品推荐

