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

如何在Airflow DAG中异步调用AWS Lambda并获取结果?

解决Airflow异步调用Lambda后获取返回结果的问题

你当前使用apache-airflow-providers-amazon 7.4.1通过LambdaInvokeFunctionOperator的invocation_type="Event"实现Lambda异步调用时,无法直接获取返回结果的核心原因是:异步调用模式下,Airflow仅会收到Lambda的请求ID(而非执行结果),且LambdaFunctionStateSensor监控的是Lambda函数本身的状态(如Active/Inactive),无法跟踪单次调用的执行状态和结果。

以下是两种可行的解决方案:

方案一:Lambda执行完将结果写入存储,Airflow用传感器等待并读取

让Lambda函数在执行完成后,将结果写入S3(或DynamoDB等存储服务),Airflow通过对应存储的传感器等待结果生成,再读取内容。

步骤1:修改Lambda函数,添加结果写入逻辑

更新你的"hello world" Lambda,将结果写入S3:

import boto3
s3 = boto3.client('s3')

def lambda_handler(event, context):
    result = "hello world"
    # 用Lambda请求ID作为文件名,避免重复
    s3.put_object(
        Bucket='your-target-s3-bucket',
        Key=f'lambda-outputs/{context.aws_request_id}.txt',
        Body=result.encode('utf-8')
    )
    return result

步骤2:修改Airflow DAG,实现等待与读取

更新你的DAG,添加请求ID提取、S3传感器等待、结果读取的逻辑:

from datetime import datetime, timedelta

from airflow.models.dag import DAG
from airflow.providers.amazon.aws.operators.lambda_function import LambdaInvokeFunctionOperator
from airflow.providers.amazon.aws.sensors.s3 import S3KeySensor
from airflow.operators.python_operator import PythonOperator
from airflow.providers.amazon.aws.hooks.s3 import S3Hook

default_args = {
    'owner': 'airflow',
    'depends_on_past': False,
    'email_on_failure': False,
    'email_on_retry': False,
}

def extract_request_id(**context):
    # 从Lambda调用任务的XCom中获取请求ID
    invoke_resp = context['ti'].xcom_pull(task_ids='setup__invoke_lambda_function')
    request_id = invoke_resp['ResponseMetadata']['RequestId']
    context['ti'].xcom_push(key='lambda_request_id', value=request_id)

def read_lambda_result(**context):
    request_id = context['ti'].xcom_pull(key='lambda_request_id')
    s3_hook = S3Hook(aws_conn_id='aws_default')
    # 读取S3中存储的Lambda结果
    result = s3_hook.read_key(
        key=f'lambda-outputs/{request_id}.txt',
        bucket_name='your-target-s3-bucket'
    )
    print(result)  # 输出"hello world"

with DAG(
    'lambda-test',
    default_args=default_args,
    description='Runs a lambda as a test',
    schedule_interval=timedelta(minutes=20),
    start_date=datetime(2021, 1, 1),
    catchup=False,
) as dag:

    invoke_lambda = LambdaInvokeFunctionOperator(
        task_id='setup__invoke_lambda_function',
        function_name="aws-pipeline-lambdas-dev-hello-world",
        invocation_type="Event"
    )

    get_request_id = PythonOperator(
        task_id='extract_lambda_request_id',
        python_callable=extract_request_id,
        provide_context=True
    )

    wait_for_result = S3KeySensor(
        task_id='wait_for_lambda_result',
        bucket_key=f'lambda-outputs/{{{{ ti.xcom_pull(key="lambda_request_id") }}}}.txt',
        bucket_name='your-target-s3-bucket',
        aws_conn_id='aws_default',
        poke_interval=10,  # 每10秒检查一次结果是否生成
        timeout=300  # 最长等待5分钟
    )

    read_result = PythonOperator(
        task_id='Read-Lambda-Output',
        python_callable=read_lambda_result,
        provide_context=True
    )

    # 任务依赖链
    invoke_lambda >> get_request_id >> wait_for_result >> read_result

方案二:配置Lambda异步调用目的地,Airflow监听消息

通过AWS Lambda控制台配置异步调用目的地,让Lambda执行完成后将结果发送到SQS/SNS/EventBridge,Airflow使用对应传感器监听消息并获取结果:

  1. 进入Lambda函数配置页面,找到"异步调用"设置,添加目的地(如SQS队列);
  2. 在Airflow DAG中使用SqsSensor等待队列中的消息,读取消息内容即为Lambda返回结果。

说明

  • 两种方案中,方案一更适合简单的结果存储场景,方案二更适合需要异步通知的复杂工作流;
  • 需确保Airflow的AWS连接权限足够访问对应存储服务或消息队列。

内容的提问来源于stack exchange,提问作者TrendSpark

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.23 10:54:59