You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

与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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.11 08:06:30