- 示例工程
【免费下载链接】python-docs-samples
Code samples used on cloud.google.com
导读
本指南围绕 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/Sub | eventarc/pubsub/ | Cloud Pub/Sub 消息 | 解码消息 payload,返回问候语并携带事件 ID |
| Audit Logs – Cloud Storage | eventarc/audit-storage/ | Cloud Storage 写操作审计日志 | 从 CloudEvent 中提取存储桶 subject 并打印 |
| Generic | eventarc/generic/ | 任意 HTTP 触发的 CloudEvent | 原样回显事件请求头与请求体 |
此外,仓库中还包含两个进阶变体示例:audit_iam(IAM 审计日志)与storage_handler(Cloud Storage CloudEvent 直连),下文将一并结合源码分析。
环境准备与本地运行
按照 README 的 Setup 章节,本地运行这些示例需要以下步骤:
- 配置 Cloud Run 开发环境:示例最终部署目标是 Cloud Run,因此需要先完成 Cloud Run 的本地开发环境配置(包括
gcloudCLI、Docker、对目标 GCP 项目的权限等)。 - 安装 Python 工具链:安装
pip与virtualenv。若从零开始,可参考 Google Cloud 官方的 Python 开发环境设置指南。 - 创建并激活虚拟环境(README 注明示例兼容 Python 2.7 与 3.4+;实际仓库依赖已更新至现代版本,见下文依赖分析):
virtualenv env source env/bin/activate - 安装依赖:
pip install -r requirements.txt - 本地启动应用:
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 为例,其构建步骤具有代表性:
- 基础镜像:
FROM python:3.14-slim,使用官方精简 Python 镜像; - 日志即时刷新:
ENV PYTHONUNBUFFERED True,保证日志立即出现在 Cloud Run 日志中; - 依赖分层缓存:先单独
COPY requirements.txt ./并RUN pip install -r requirements.txt,避免每次代码变更都重新安装依赖; - 工作目录:
ENV APP_HOME /app与WORKDIR $APP_HOME; - 启动命令(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
相关推荐
从零开始构建智能数字人助理:Fay框架如何让AI真正为你工作
从零开始构建智能数字人助理:Fay框架如何让AI真正为你工作 想象一下,每天早上醒来,你的数字人助理已经为你规划好了一天的工作日程;当你需要查询信息时,它不仅能
Nikto与Google Cloud Functions Eventarc集成:事件驱动扫描新范式
Nikto与Google Cloud Functions Eventarc集成:事件驱动扫描新范式 痛点与解决方案 你是否还在为Web服务器安全检测的及时性和自
rembg 人像分割实操指南:从发丝边缘到批量证件照
rembg 人像分割实操指南:从发丝边缘到批量证件照 这篇文章讲的是 rembg 的人像背景去除这条线:为什么默认模型抠人像总差点意思、birefnet por
人工智能计算机视觉图像处理
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考