如何在Google Cloud Composer流水线中调用Cloud Function?
报错根因
你遇到的JSONDecodeError是在解析Cloud Function返回的响应时触发的:你在response_check回调中直接调用了response.json(),但Cloud Function返回的内容为空,或者格式不是合法JSON,导致解析失败。
常见触发场景有三个:
- Cloud Function执行成功后没有返回JSON结构,只返回了空字符串、204无内容状态码,或者返回了纯文本、HTML内容
- 调用Cloud Function时鉴权失败,GCP返回了非JSON格式的错误页面
- SimpleHTTPOperator没有携带正确的身份认证头,请求被拦截返回非JSON响应
修复方案
方案1:直接修复现有SimpleHTTPOperator的逻辑
- 先调整响应检查逻辑,先判断状态码和响应内容是否为空,再做JSON解析:
def custom_response_check(response): # 先打印原始响应方便排查,调试完成后可以删除 print(f"响应状态码: {response.status_code}") print(f"响应原始内容: {response.text}") # 合法状态码校验 if response.status_code not in [200, 201, 204]: return False # 空响应直接返回成功 if len(response.content.strip()) == 0: return True # 非空响应再解析JSON return len(response.json()) > 0 # 替换你原来的lambda表达式 response_check=custom_response_check
- 调整Cloud Function的返回逻辑:如果需要返回数据,明确返回JSON结构和200状态码;如果不需要返回数据,直接返回204状态码,不要返回空字符串:
# 需要返回数据的场景 def cf_handler(request): # 你的业务逻辑 return {"status": "success", "data": "xxx"}, 200 # 不需要返回数据的场景 def cf_handler(request): # 你的业务逻辑 return "", 204
- 给Cloud Composer运行时的服务账号授予
Cloud Functions Invoker角色,确保有权限调用目标Cloud Function,避免鉴权失败。
方案2:使用官方Cloud Function调用Operator(更推荐)
Airflow的Google官方Provider包已经提供了专门调用Cloud Function的Operator,不需要自己处理鉴权、请求头、响应解析等逻辑,比SimpleHTTPOperator适配性更好:
from airflow.providers.google.cloud.operators.functions import CloudFunctionInvokeFunctionOperator invoke_cf_task = CloudFunctionInvokeFunctionOperator( task_id="invoke_my_cf", project_id="你的GCP项目ID", location="Cloud Function部署的区域,例如us-central1", function_id="你的Cloud Function名称", input_data={"param1": "val1", "param2": "val2"}, # 传给CF的参数,不需要可以省略 gcp_conn_id="google_cloud_default", # Composer默认已经配置好该连接,无需额外修改 )
内容的提问来源于stack exchange,提问作者Snehil Singh
相关产品推荐
相关产品推荐

