拓十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Eventarc Python 示例:基于 Cloud Run 的事件驱动接收器开发实战

Eventarc Python 示例:基于 Cloud Run 的事件驱动接收器开发实战
  • 示例工程

【免费下载链接】python-docs-samples

Code samples used on cloud.google.com

项目地址:https://gitcode.com/GitHub_Trending/py/python-docs-samples
点击查看免费下载

导读

本指南围绕 eventarc/README.md 展开,介绍如何用 Flask 编写接收 CloudEvents 事件的服务端应用,并部署到 Cloud Run,用于响应 Pub/Sub 消息、Cloud Storage 审计日志等事件。读完本文,你将掌握 Pub/Sub、审计日志(Audit Logs)与通用(Generic)三类事件接收器的核心实现、本地运行、自动化测试与容器化部署的完整实战方案。

Eventarc 是什么

Eventarc 是 Google Cloud 上统一的事件路由服务,它将来自不同事件源(如 Cloud Pub/Sub、Cloud Storage、IAM 审计日志等)的事件统一封装为 CloudEvents 格式,并投递到事件接收端点(通常是一个 Cloud Run 服务)。README 明确说明该目录包含用于 Eventarc 的 Python 示例,这些示例以 Flask 为基础,聚焦于"接收事件并做出响应"这一核心职责。

示例全景

README 将本目录的示例划分为三类,对应 Eventarc 最常见的三种接入场景:

示例目录(仓库根路径)事件类型核心处理逻辑
Pub/Subeventarc/pubsub/Cloud Pub/Sub 消息解码消息 payload,返回问候语并携带事件 ID
Audit Logs – Cloud Storageeventarc/audit-storage/Cloud Storage 写操作审计日志从 CloudEvent 中提取存储桶 subject 并打印
Genericeventarc/generic/任意 HTTP 触发的 CloudEvent原样回显事件请求头与请求体

此外,仓库中还包含两个进阶变体示例:audit_iam(IAM 审计日志)与storage_handler(Cloud Storage CloudEvent 直连),下文将一并结合源码分析。

环境准备与本地运行

按照 README 的 Setup 章节,本地运行这些示例需要以下步骤:

  1. 配置 Cloud Run 开发环境:示例最终部署目标是 Cloud Run,因此需要先完成 Cloud Run 的本地开发环境配置(包括gcloudCLI、Docker、对目标 GCP 项目的权限等)。
  2. 安装 Python 工具链:安装pip与virtualenv。若从零开始,可参考 Google Cloud 官方的 Python 开发环境设置指南。
  3. 创建并激活虚拟环境(README 注明示例兼容 Python 2.7 与 3.4+;实际仓库依赖已更新至现代版本,见下文依赖分析):
    virtualenv env source env/bin/activate
  4. 安装依赖:
    pip install -r requirements.txt
  5. 本地启动应用:
    python main.py

启动后,Flask 开发服务器默认监听0.0.0.0:8080(可通过PORT环境变量覆盖),事件源或测试工具即可向/端点发送 POST 请求。

依赖与运行配置分析

从仓库实际的依赖清单可以确认各示例的运行环境:

  • eventarc/pubsub/requirements.txt:Flask==3.1.3、gunicorn==23.0.0
  • eventarc/audit-storage/requirements.txt:在上述基础上增加cloudevents==1.11.0(CloudEvents 解析库)
  • eventarc/audit_iam/requirements.txt:增加google-events==0.14.0(Google 事件 protobuf 类型)、cloudevents==1.11.0、googleapis-common-protos==1.66.0
  • eventarc/storage_handler/requirements.txt:与 audit_iam 相同,依赖google-events与cloudevents

可见,处理结构化事件(审计日志、GCS 对象事件)的示例需要额外的cloudevents与google-events依赖;而 Pub/Sub 与 Generic 示例仅需 Flask 即可。所有示例在main.py中统一使用如下启动方式:

if __name__ == "__main__": app.run(debug=True, host="0.0.0.0", port=int(os.environ.get("PORT", 8080)))

示例一:Pub/Sub 消息接收器

核心实现

eventarc/pubsub/main.py 定义了一个仅接受POST方法的 Flask 路由。Eventarc 投递的 Pub/Sub 消息以 JSON 结构到达,关键字段为message.data——该字段是 Pub/Sub 消息 payload 的base64 编码字符串,因此处理前必须先解码:

@app.route("/", methods=["POST"]) def index(): data = request.get_json() if not data: msg = "no Pub/Sub message received" print(f"error: {msg}") return f"Bad Request: {msg}", 400 if not isinstance(data, dict) or "message" not in data: msg = "invalid Pub/Sub message format" print(f"error: {msg}") return f"Bad Request: {msg}", 400 pubsub_message = data["message"] name = "World" if isinstance(pubsub_message, dict) and "data" in pubsub_message: name = base64.b64decode(pubsub_message["data"]).decode("utf-8").strip() resp = f"Hello, {name}! ID: {request.headers.get('ce-id')}" print(resp) return (resp, 200)

实现要点:

  • 空 payload 防护:请求体为空或非 JSON 时返回400,避免服务异常;
  • 结构校验:要求顶层必须是 dict 且包含message键,否则返回400;
  • base64 解码:pubsub_message["data"]需经base64.b64decode(...).decode("utf-8")还原为原始字符串,再.strip()去除空白;
  • CloudEvent ID 回显:通过request.headers.get('ce-id')读取 CloudEvent 标准头,将事件 ID 拼入响应,便于链路追踪。

测试用例印证

eventarc/pubsub/main_test.py 使用 Flask 内置的 test client 覆盖了五类场景,与源码的校验逻辑一一对应:

  • test_empty_payload:空 payload 返回400;
  • test_invalid_payload:缺少message键返回400;
  • test_invalid_mimetype:非 JSON 内容返回400;
  • test_minimally_valid_message:仅{"message": true}返回200,并打印Hello, World! ID: {ce-id};
  • test_populated_message:构造 base64 编码的用户名,验证解码结果正确出现在响应与日志中。

测试通过binary_headers模拟 CloudEvent 二进制头(ce-id、ce-type、ce-source、ce-specversion),这与 Eventarc 实际投递时携带的 CloudEvent 头一致。

示例二:Audit Logs – Cloud Storage 审计日志接收器

核心实现

eventarc/audit-storage/main.py 演示了如何接收由 Cloud Storage 写操作产生的审计日志事件。与 Pub/Sub 示例不同,这里引入cloudevents.http.from_http将原始 HTTP 请求转换为标准 CloudEvent 对象:

from cloudevents.http import from_http from flask import Flask, request app = Flask(__name__) @app.route("/", methods=["POST"]) def index(): # Create a CloudEvent object from the incoming request event = from_http(request.headers, request.data) # Gets the GCS bucket name from the CloudEvent # Example: "storage.googleapis.com/projects/_/buckets/my-bucket" bucket = event.get("subject") print(f"Detected change in Cloud Storage bucket: {bucket}") return (f"Detected change in Cloud Storage bucket: {bucket}", 200)

实现要点:

  • CloudEvent 解析:from_http(request.headers, request.data)自动从请求头与 body 重建 CloudEvent,无需手写解析;
  • subject 语义:CloudEvent 的subject属性携带事件主体信息。对于 Cloud Storage 审计日志,其值为存储桶的完整资源名,形如storage.googleapis.com/projects/_/buckets/my-bucket,源码注释给出了这一格式示例;
  • 业务响应:将检测到的存储桶名拼入日志与 HTTP 响应,返回200。

测试用例印证

eventarc/audit-storage/main_test.py 演示了 CloudEvents 库的"反向"用法——先用cloudevents.http.CloudEvent构造事件对象,再用to_binary()序列化为二进制 HTTP 请求,最后通过 test client 投递:

ce_attributes = { "id": str(uuid4), "type": "com.pytest.sample.event", "source": "<my-test-source>", "specversion": "1.0", "subject": "test-bucket", } event = CloudEvent(ce_attributes, dict()) headers, body = to_binary(event) r = client.post("/", headers=headers, data=body) assert "Detected change in Cloud Storage bucket: test-bucket" in r.text

测试断言subject属性被正确读取为存储桶名,验证了event.get("subject")的取值链路。

示例三:Generic 通用事件接收器

核心实现

eventarc/generic/main.py 是一个"回显"型接收器,用于调试与观察任意 Eventarc 事件的全貌。它打印事件请求头与请求体,并将其原样返回:

@app.route("/", methods=["POST"]) def index(): print("Event received!") print("HEADERS:") headers = dict(request.headers) headers.pop("Authorization", None) # do not log authorization header if exists print(headers) print("BODY:") body = dict(request.json) print(body) resp = {"headers": headers, "body": body} return (resp, 200)

实现要点:

  • 全量回显:将请求头与请求体同时写入日志并作为 JSON 响应返回,便于开发者直观核对事件内容;
  • 安全处理:显式移除Authorization头,避免将认证凭据打印进日志——这是事件接收器普遍需要遵守的安全实践;
  • 适用场景:当接入一种新的事件源、需要确认事件结构与字段名时,Generic 接收器是最快的诊断工具。

测试用例印证

eventarc/generic/main_test.py 验证了回显逻辑的完整性:

  • 断言日志中出现Event received!;
  • 断言响应体中的头包含"Ce-Specversion":"1.0"(验证 CloudEvent 头被回显);
  • 断言请求体内容{'message': {'data': 'Hello'}}被完整打印。

进阶变体:审计日志的结构化处理

仓库中还提供了两个超出 README 列表演示的进阶示例,它们展示了同一主题下更深入的生产级用法。

IAM 审计日志(audit_iam)

eventarc/audit_iam/main.py 将审计日志事件反序列化为强类型 protobuf 对象google.events.cloud.audit.LogEntryData,实现"检测服务账号密钥创建"的业务逻辑:

event = from_http(request.headers, request.get_data()) log_entry = LogEntryData.from_json( json.dumps(event.get_data()), ignore_unknown_fields=True ) if log_entry.proto_payload.service_name != "iam.googleapis.com": return ("Received event was not from IAM.", 400) if log_entry.proto_payload.status.code != 0: return ("Key creation failed, not reporting.", 204) user = log_entry.proto_payload.authentication_info.principal_email service_account = log_entry.proto_payload.request["name"] keypath = log_entry.proto_payload.response["name"] print(f"New Service Account Key created for {service_account} by {user}: {keypath}")

关键细节:

  • 大小写转换:LogEntryData.from_json内部将 JSON 风格的lowerCamelCase字段名转换为 protobuf 风格的snake_case;
  • 未知字段容忍:ignore_unknown_fields=True用于跳过审计日志中的@type等元数据字段;
  • 多级过滤:先校验服务名是否为iam.googleapis.com(否则400),再校验操作状态码code != 0表示失败(返回204不告警);
  • 可直接操作的输出:注释指出keypath可用于gcloud iam service-accounts keys disable ${keypath}等后续处置。

Cloud Storage CloudEvent(storage_handler)

eventarc/storage_handler/main.py 展示的是 Cloud Storage 对象级事件(而非审计日志)的处理方式,使用google.events.cloud.storage.StorageObjectData强类型解析事件数据:

try: storage_obj = StorageObjectData(event.data) gcs_object = os.path.join(storage_obj.bucket, storage_obj.name) update_time = storage_obj.updated return ( f"Cloud Storage object changed: {gcs_object}" + f" updated at {update_time}", 200, ) except ValueError as e: return (f"Failed to parse event data: {e}", 400)

与审计日志示例相比,这里直接从event.data构造StorageObjectData,能够访问bucket、name、updated等对象属性;解析失败(ValueError)时返回400。eventarc/storage_handler/main_test.py 同时验证了正常路径(Cloud Storage object changed: test-bucket/my-file.txt)与异常路径(非法数据返回400)。

容器化部署:Docker 与 Cloud Run

所有示例都提供 Dockerfile,用于将 Flask 应用打包为 Cloud Run 镜像。以 eventarc/audit-storage/Dockerfile 为例,其构建步骤具有代表性:

  1. 基础镜像:FROM python:3.14-slim,使用官方精简 Python 镜像;
  2. 日志即时刷新:ENV PYTHONUNBUFFERED True,保证日志立即出现在 Cloud Run 日志中;
  3. 依赖分层缓存:先单独COPY requirements.txt ./并RUN pip install -r requirements.txt,避免每次代码变更都重新安装依赖;
  4. 工作目录:ENV APP_HOME /app与WORKDIR $APP_HOME;
  5. 启动命令(gunicorn 生产服务器):
    CMD exec gunicorn --bind :$PORT --workers 1 --threads 8 --timeout 0 main:app

Dockerfile 注释明确解释了该命令的设计意图:单 worker + 8 线程;若部署环境具有多 CPU 核心,应将--workers调整为与核心数一致。--timeout 0用于禁用请求超时,适配事件处理可能耗时较长的场景。$PORT由 Cloud Run 注入。eventarc/storage_handler/Procfile 则为非容器平台提供了等价的web: gunicorn --bind :$PORT --workers 1 --threads 8 --timeout 0 main:app启动声明。

本地测试与验证

每个示例均包含main_test.py,测试框架统一使用 pytest(见各目录的requirements-test.txt,如 eventarc/pubsub/requirements-test.txt)。运行方式:

pip install -r requirements.txt -r requirements-test.txt pytest

测试的共性模式可以总结为三类:

  • 模拟 CloudEvent 头:通过字典构造ce-id、ce-type、ce-source、ce-specversion等标准头;
  • 构造事件对象:使用cloudevents.http.CloudEvent与cloudevents.conversion.to_binary将事件序列化为真实 HTTP 请求(见 eventarc/audit-storage/main_test.py);
  • 强类型数据构造:对于依赖google-events的示例,直接用StorageObjectData(bucket=..., name=...)、StorageObjectData.to_dict()构造合法负载(见 eventarc/storage_handler/main_test.py)。

总结

从 README 出发,结合仓库源码可以梳理出 Eventarc Python 事件接收器的完整开发脉络:

  • 三类核心模式:Pub/Sub 消息的 base64 解码、审计日志的 CloudEvent 解析(subject属性)、Generic 的全量回显,覆盖了事件接收的绝大多数场景;
  • 两种解析路线:轻量场景直接用request.get_json()手写解析;生产级场景用cloudevents+google-events强类型解析(LogEntryData、StorageObjectData),并借助ignore_unknown_fields、多级状态码过滤提升健壮性;
  • 一套部署范式:Flask 应用 + gunicorn(单 worker 多线程、--timeout 0)+ Dockerfile 分层缓存,是 Cloud Run 事件接收服务的事实标准配置;
  • 测试方法论:以 pytest + Flask test client 模拟 CloudEvent 二进制请求,兼顾正常路径与400/204异常路径的断言。

无论你是要接入第一条 Pub/Sub 事件,还是需要处理结构复杂的 IAM 审计日志,本目录的示例都可以直接作为脚手架,结合 README 的 Setup 流程快速落地。

  • 示例工程

【免费下载链接】python-docs-samples

Code samples used on cloud.google.com

项目地址:https://gitcode.com/GitHub_Trending/py/python-docs-samples
点击查看免费下载

相关推荐

上一篇:为什么选择nxdumptool:3个超越传统备份方案的关键优势
下一篇:ChatGPT Retrieval Plugin 接入 Postgres + pgvector:环境配置、迁移与索引优化实战指南

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

返回列表