与AWS S3桶交互时Airflow DAG卡在运行状态的问题求助
Airflow DAG与S3交互时挂起,单独运行代码正常
简单DAG可正常执行,但涉及AWS S3桶交互的DAG会卡在运行状态。使用Airflow S3 Hook实现「调用API获取JSON并保存至S3」的功能,单独运行函数时工作正常,嵌入DAG执行则出现挂起。可替换桶名进行测试。
功能函数代码
from airflow.providers.amazon.aws.hooks.s3 import S3Hook import requests import json def dump_json_to_s3(url, bucket_name, file_name): """ :param url: API地址 :param bucket_name: S3桶名 :param file_name: JSON文件名 :return: """ try: print("开始处理...") # 从API获取JSON数据 response = requests.get(url=url) # S3键路径 key_path = f"asset/{file_name}.json" # 初始化S3 Hook hook = S3Hook('s3_conn') print(response.text) hook.load_string( string_data=response.text, key=key_path, bucket_name=bucket_name, replace=True ) # 打印成功日志 print(f"成功:JSON文件已上传至S3路径 {key_path}") # 异常捕获 except requests.exceptions.RequestException as e: print(f"请求错误:{e}") except boto3.exceptions.S3UploadFailedError as e: print(f"S3上传错误:{e}") except Exception as e: print(f"未知错误:{e}") return # 单独测试函数可执行以下代码 url = "http://engineering-exam.s3-website.ap-southeast-2.amazonaws.com/" bucket_name = "替换为你的S3桶名" dump_json_to_s3(url=url, bucket_name=bucket_name, file_name='asset_new')
DAG代码
from airflow import DAG from airflow.operators.python import PythonOperator from datetime import datetime, timedelta from dump_json_to_s3_file import * default_args = { 'owner': 'airflow', 'start_date': datetime(2023, 9, 1), } dag = DAG( dag_id='dump_to_s3', schedule_interval=None, default_args=default_args, catchup=False ) convert_json = PythonOperator( task_id="convertor", python_callable=dump_json_to_s3, # 引用函数而非直接调用 op_args=[ # 传入参数列表 "http://engineering-exam.s3-website.ap-southeast-2.amazonaws.com/", "替换为你的S3桶名", "airflow_tester" ], dag=dag ) convert_json
可能的原因及解决办法
- Airflow连接配置错误:确认Airflow中
s3_conn连接的配置正确性。检查连接类型为Amazon Web Services,Access Key/Secret Key具备S3读写权限,区域与目标桶匹配;若使用IAM角色(如EC2实例角色),需在连接的Extra字段配置{"role_arn": "arn:aws:iam::xxx:role/xxx"},且角色信任关系包含Airflow Worker的身份。 - 网络与权限限制:检查Airflow Worker所在环境(容器/EC2等)是否允许出站访问S3,安全组、NACL规则是否放行对应流量;同时确认目标S3桶的政策允许Airflow使用的身份执行
PutObject操作。 - 超时与资源配置:初始化S3 Hook时添加超时参数,避免请求无限等待:
hook = S3Hook('s3_conn', config_kwargs={"connect_timeout": 10, "read_timeout": 30});在PythonOperator中配置execution_timeout=timedelta(minutes=5),强制超时失败以便排查。 - 执行环境不一致:对比单独运行代码的环境与Airflow Worker的环境,确保Python版本、
boto3、airflow-providers-amazon等依赖包版本一致,版本不兼容可能导致隐性挂起。 - 查看详细日志:不要仅依赖UI状态,通过命令
airflow logs -t convertor -d dump_to_s3查看Worker端实时日志,获取未捕获的异常或权限提示信息。
内容的提问来源于stack exchange,提问作者d789w
相关产品推荐
相关产品推荐

