Airflow 2.7.3中LambdaInvokeFunctionOperator传递XCom生成字节数组Payload问题
解决Airflow 2.7.3中LambdaInvokeFunctionOperator使用XCom生成字节Payload的问题
方案1:模板字符串+自动字节编码
LambdaInvokeFunctionOperator的payload参数支持Jinja模板解析,且Airflow会自动将解析后的字符串转换为字节数组。直接把JSON结构写成带Jinja引用的字符串模板即可:
from airflow.providers.amazon.aws.operators.lambda_function import LambdaInvokeFunctionOperator invoke_lambda = LambdaInvokeFunctionOperator( task_id="invoke_lambda_task", function_name="your-lambda-function-name", payload='{"schema": "{{ ti.xcom_pull(task_ids=\'start-process\', key=\'schema\') }}", "ids": {{ ti.xcom_pull(task_ids=\'start-process\', key=\'ids\') }}}', aws_conn_id="aws_default" )
注意事项:
ids是整数列表,Jinja解析后直接输出原生列表结构,无需加引号schema是字符串,模板变量需要用双引号包裹- 整个payload是JSON格式的字符串模板,Airflow完成模板解析后会自动转为字节数组传给Lambda
方案2:Python函数生成字节Payload(复杂场景适用)
如果需要对XCom数据做额外处理(比如过滤、格式转换),可以先用PythonOperator生成字节数组再通过XCom传递:
from airflow.operators.python import PythonOperator from airflow.providers.amazon.aws.operators.lambda_function import LambdaInvokeFunctionOperator import json def generate_lambda_payload(**context): ti = context["ti"] schema = ti.xcom_pull(task_ids="start-process", key="schema") ids = ti.xcom_pull(task_ids="start-process", key="ids") payload_dict = {"schema": schema, "ids": ids} # 转为字节数组并存入XCom ti.xcom_push(key="lambda_payload", value=json.dumps(payload_dict).encode("utf-8")) generate_payload_task = PythonOperator( task_id="generate_lambda_payload", python_callable=generate_lambda_payload, provide_context=True ) invoke_lambda = LambdaInvokeFunctionOperator( task_id="invoke_lambda_task", function_name="your-lambda-function-name", payload="{{ ti.xcom_pull(task_ids=\'generate_lambda_payload\', key=\'lambda_payload\') }}", aws_conn_id="aws_default" ) generate_payload_task >> invoke_lambda
核心避坑点
- 不要直接用
bytearray(json.dumps(...))包裹Jinja模板,这会导致模板被当作普通字符串,无法触发XCom解析 - 确保XCom中
ids的类型是列表,Jinja解析时会直接输出符合JSON格式的列表结构
内容的提问来源于stack exchange,提问作者Kari Marttila
相关产品推荐
相关产品推荐

