Google PDE: 消息传递、事件注入与实时服务 — 学习指南
属于 Google Professional Data Engineer — 学习指南. 使用经过验证的答案练习: Google 考试中心, 或参加限时模拟考试: ExamRoll.io.
概览
Google Cloud 上的消息传递、事件注入和实时服务以 Cloud Pub/Sub 和 Eventarc 为核心,提供解耦、持久化的传输;以 Dataflow 提供有状态的流处理;并以 BigQuery、Cloud Storage 和各种运维型数据库作为数据汇 (sink)。通过为至少一次 (at-least-once) 投递、幂等性消费和可观测性进行设计,可以确保系统具有弹性,能够在故障、背压和模式演进的情况下弹性伸缩,同时保持正确性。
使用 Pub/Sub 实现核心消息传递
- 主题和订阅
- 发布者将消息发送到主题;订阅者通过订阅进行附加(多个订阅者可以独立消费相同的消息)。
- 订阅类型:
- 拉取 (Pull):客户端显式拉取消息;使用流式拉取 (streaming pull) 可获得最高吞吐量并减少网络往返次数。
- 推送 (Push):Pub/Sub 通过 HTTPS 投递消息;您的端点必须返回 2xx 状态码以示确认。
- 导出到 BigQuery:BigQuery 订阅无需代码即可将消息投递到 BigQuery 表;当载荷 (payload) 与声明的模式匹配且需要低延迟注入分析系统时,此为最佳选择。
- 排序键 (Ordering keys)
- 在主题和订阅上启用消息排序,以实现按排序键的有序投递。每个键的吞吐量是串行化的:每个键只有一个在途 (in-flight) 消息,这可能会阻塞后续消息;使用大量键(例如,hash(device_id))来进行扩展。
- 扇出 (Fan-out) 和重放 (replay)
- 为不同的消费者创建独立的订阅,以隔离工作负载和保留策略。
- 使用寻址 (Seek) 或快照 (snapshot) 从指定时间戳或快照进行重放,以用于恢复和数据回填。
权衡取舍:
- 排序会降低每个键的并行度和吞吐量;除非严格需要,否则应禁用排序。
- 推送简化了客户端代码,但引入了 HTTP 端点扩展、安全性和退避 (backoff) 等问题;拉取则在实现高吞吐量时提供了更多的控制和稳定性。
投递语义、确认、保留和死信
- 确认 (Acknowledgment) 和截止时间 (deadlines)
- 至少一次 (At-least-once) 投递:可能会出现重复消息。
- 每次投递都有一个确认截止时间 (ack deadline,默认为 10 秒)。在处理长时间运行的工作时,可以延长该时间 (ModifyAckDeadline);未在截止时间前确认是导致推送投递重复的最常见原因。
- 否定确认 (Nack) 或截止时间到期会使消息有资格被重新投递。
- 保留 (Retention)
- 未确认的消息会在订阅的确认截止时间内被保留并重试;已确认的消息最多可以被保留主题所设定的消息保留时长,以供重放。配置保留时长,使其能覆盖您的最大故障时长加上恢复时间。
- 重试 (Retries)
- 拉取 (Pull):在确认截止时间到期后会进行重新投递;使用流控制 (flow control) 限制来控制并发。
- 推送 (Push):采用指数退避 (exponential backoff) 策略;只有 HTTP 2xx 状态码表示成功。3xx/4xx/5xx 会触发重试。应实现幂等处理程序以容忍重复投递。
- 死信主题 (Dead-letter topics, DLTs)
- 为每个订阅配置一个死信主题和最大投递尝试次数,以隔离“毒丸”消息 (poison messages)。
- 监控死信队列 (DLQ) 的消息量;创建分类处理工作流,并在修正后将消息重新发布到主主题。
示例:
undefined
投递语义总结:
- Pub/Sub:至少一次投递,如果启用了排序键,则在排序键内尽力而为 (best-effort) 地保证顺序。
- 数据汇 (Sinks):BigQuery 的插入 API 提供了重复数据缓解机制(如 insertId 或 Storage Write API 的流偏移量),但消费者和写入者仍应设计为幂等的。
模式、兼容性和验证
- Pub/Sub 模式
- 原生支持 Avro 和 Protocol Buffers,其模式可被集中存储。
- 主题级别的模式设置:编码 (Avro 或 Protobuf) 和强制执行策略 (无、仅验证或要求强制)。
- 生产者发布编码后的载荷;当启用强制执行策略时,Pub/Sub 会根据当前模式对消息进行验证。
- 演进和兼容性
- 使用向后兼容的变更(例如,添加可选字段,在 Avro 中添加带默认值的字段,在 Protobuf 中永不重用标签,避免移除或重命名字段)。
- 明确地对模式进行版本控制。对于破坏性变更,可以双重发布到 v1 和 v2 两个主题,或者添加一个版本字段并据此进行路由。
- 生产者-消费者契约
- 消费者应该忽略未知字段,并为缺失的字段设置默认值。
- 在发布到生产环境前,测试模式在所有消费者中的兼容性;在预演 (staging) 环境的订阅中进行验证,并使用与生产环境相同的模式强制执行策略。
Avro 简短示例(摘录):
undefined
事件驱动集成、Eventarc 和 Kafka 互操作性
- Eventarc 和 CloudEvents
- Eventarc 使用 CloudEvents 规范将来自 Google Cloud 服务(以及通过 Pub/Sub 的自定义来源)的事件路由到 Cloud Run、GKE 或 Workflows。诸如 type、source、subject 等属性可实现精细的过滤和可审计性。
- 使用属性过滤器来最小化扇出(fan-out)并减少下游负载。
- 交付语义为至少一次(at-least-once);应尽可能使处理程序(handler)具有幂等性和无状态性。
- Eventarc 触发器示例:
gcloud eventarc triggers create gcs-finalize-to-run
–destination-run-service=ingestor
–event-filters=“type=google.cloud.storage.object.v1.finalized”
–event-filters=“bucket=my-data-bucket”
–service-account=eventarc-sa@PROJECT_ID.iam.gserviceaccount.com - Kafka 互操作性与托管式迁移
- Dataflow 模板可连接 Kafka <-> Pub/Sub 以进行分阶段迁移。镜像主题并保留键(key);先切换消费者,然后是生产者,或在过渡期间进行双写。
- Pub/Sub Lite 提供分区化、容量预配的流式处理,具有基于键的路由和更低的成本;它是区域级/可用区级的,适用于类似 Kafka 的工作负载,其主要关注点是可预测的容量和每个分区的顺序。
- 迁移注意事项:
- 顺序:将 Kafka 的键(key)映射到 Pub/Sub 的排序键(ordering key)或 Lite 的分区。
- 偏移量(Offset):将偏移量作为消息属性携带以用于诊断;迁移后,消费者不能再依赖 Kafka 的偏移量。
- 交付:接受至少一次(at-least-once)的交付语义;在下游强制实现幂等性。
- 模式(Schema):将 Confluent Schema Registry 的定义迁移到 Pub/Sub 模式,或在 Protobuf/Avro 上进行标准化,并采用兼容的演进规则。
流式摄取模式、吞吐量、扩缩、安全性和运维
- 实时摄取模式
- Pub/Sub -> Dataflow -> BigQuery:使用 BigQuery Storage Write API sink 以实现高吞- 吐量和基于流偏移量的幂等性;将失败的记录路由到死信表以便检查。
- Pub/Sub -> Dataflow -> Cloud Storage:归档原始事件以供再处理;使用窗口化、压缩写入来平衡成本和延迟。
- Pub/Sub -> 运营存储:写入 Bigtable 以实现低延迟查找,写入 Spanner 以实现强一致性事务,或根据工作负载需求写入 Cloud SQL/Firestore。确保基于唯一事件 ID 的幂等性更新插入(upsert)。
- 至少一次(At-least-once)、重复数据预防和幂等性
- 在每条消息中携带唯一的
event_id和event_time;强制在生产者端生成 UUID。 - BigQuery 流式去重:设置
insertId或使用带有有序流的 Storage Write API;但查询时仍需使用去重逻辑进行保护。 - 查询时去重示例: WITH ranked AS ( SELECT t.*, ROW_NUMBER() OVER (PARTITION BY event_id ORDER BY event_time DESC) AS rn FROM dataset.events t ) SELECT * EXCEPT(rn) FROM ranked WHERE rn = 1;
- 对于推送端点,仅在成功处理后返回 2xx 状态码;否则应预期消息会被重传。
- 在每条消息中携带唯一的
- 消息吞吐量、配额和扩缩
- 发布者:批量发送消息并复用连接;通过多个客户端进行并行化。使用多个排序键(ordering key)来扩展有序工作负载。
- 订阅者:优先使用带流控(最大未完成字节数/消息数)的流式拉取(streaming pull)。根据处理时间调整确认截止时间(ack deadline),并在需要时延长。
- 随着数据量的增长,监控并申请提高发布和订阅吞吐量的配额;设计时应留有余量(例如,预期峰值的 2 倍)以吸收突发流量。
- 一致性和可用性
- BigQuery 流式摄取对于查询可见性是最终一致的;对于必须包含流式插入行的交互式查询,可以根据观察到的延迟(例如,P50 可用性延迟的 2 倍)进行等待,或者在 Dataflow 中使用与水印对齐的聚合进行设计,并查询物化后的结果。
- 安全性
- IAM:授予最小权限角色(为主题上的生产者授予
pubsub.publisher;为订阅上的消费者授予pubsub.subscriber)。为每个工作负载使用专用的服务账号。 - 推送认证:配置推送订阅,使其附加来自服务账号的 OIDC 令牌;在端点上强制执行受众(audience)验证。优先使用 Cloud Run 私有端点,以利用其内置的认证和 TLS。
- 加密:Pub/Sub 对传输中和静态数据进行加密;在主题上使用 CMEK 以实现客户管理的密钥。应用 VPC Service Controls 以降低数据泄露风险。如果需要,对敏感的负载字段使用客户端加密。
- IAM:授予最小权限角色(为主题上的生产者授予
- 延迟、重传和订阅者故障的运维诊断
- 使用 Cloud Monitoring 进行监控:
- 使用
subscription/num_undelivered_messages和oldest_unacked_message_age监控积压情况。 - 使用
expired_ack_deadline_count检测因确认(ack)失败而导致的重复消息。 - 使用
publish_request_count和pull_request_count监控吞吐量。
- 使用
- 要调查仪表板上缺失的事件,可以通过管道重放一个已知的数据集,并逐个阶段比较输出,以隔离出有问题的转换或 sink。
- 对于 Dataflow 流处理:
- 使用自动扩缩,并设置一个合适的
maxWorkers,以吸收来自多个源的负载。 - 对于不兼容的更新,请排空(drain)管道,以允许处理中的工作完成并防止数据丢失。
- 使用自动扩缩,并设置一个合适的
- 对于 BigQuery 插入通知,通过一个 sink 将 Cloud Logging 审计条目(过滤到特定表)路由到一个 Pub/Sub 主题以进行告警。
- 使用 Cloud Monitoring 进行监控:
实际问题场景
Contoso Freight 公司需要一个全球性的实时事件平台,用于从卡车每分钟摄取 10,000 条物联网遥测消息,丰富事件,支持交互式分析,并针对外部合作伙伴的文件上传触发工作流。一些合作伙伴的 CSV 文件包含格式错误的行,分析团队必须能够在不阻塞流的情况下检查这些错误。
- 创建核心消息传递和 schema 层
- 操作:为遥测数据定义一个 Avro schema,并将其附加到一个 Pub/Sub 主题
telemetry上,同时将 schema 强制执行设置为require。启用消息排序,并使用ordering_key = hash(device_id)进行发布。 - 理由:主题级别的 schema 强制执行能及早拒绝格式错误的事件。按设备排序可在需要时支持有序处理,而哈希能分散排序键以保持吞吐量。
- 配置具有隔离和死信功能的订阅
- 操作:为 Dataflow 创建一个拉取订阅
telemetry-stream-sub,并配置一个死信主题telemetry-dlt和max_delivery_attempts=10。添加一个 BigQuery 订阅telemetry-raw-bq,将原始事件存入一个按时间分区的表中,以用于数据溯源和重放。 - 理由:DLQ(死信队列)隔离了毒丸消息以便调查。一个独立的 BigQuery 订阅为原始事件的保存提供了一个低运维成本的导出路径,独立于处理管道。
- 构建用于数据丰富和写入 sink 的 Dataflow 流处理管道
- 操作:使用带流控的流式拉取从
telemetry-stream-sub摄取数据。根据 schema 进行验证,用参考数据丰富事件,并计算窗口聚合。使用 Storage Write API 写入 BigQuery,并指定一个命名的流和insertId = event_id;每小时将原始备份写入 Cloud Storage;将错误/失败的记录重定向到一个死信 BigQuery 表。 - 理由:Storage Write API 通过
insertId/流偏移量实现了高吞吐量、低延迟的写入和幂等性。死信表支持在不阻塞流的情况下进行检查,而 Cloud Storage 归档则支持数据重放。
- 在分析中处理重复数据和最终一致性
- 操作:对于必须排除重复项的交互式查询,在每条记录中发布
event_id和event_time,并使用一个去重视图: CREATE OR REPLACE VIEW analytics.latest_events AS SELECT * EXCEPT(rn) FROM ( SELECT e.*, ROW_NUMBER() OVER (PARTITION BY event_id ORDER BY event_time DESC) rn FROM analytics.events e ) WHERE rn = 1; 根据观察到的 BigQuery 流式可用性(例如,中位延迟的两倍)引入一个短暂的查询延迟。 - 理由:至少一次的交付语义要求幂等写入和查询时去重。考虑到流式可见性的延迟,等待可以减少因处理中行而导致的遗漏。
- 使用 Eventarc 集成合作伙伴的文件上传
- 操作:配置 Eventarc,将存储桶
partner-drops的object.finalized事件路由到一个 Cloud Run 服务。该服务启动一个批处理 Dataflow 作业来将 CSV 加载到 BigQuery,并将解析错误发送到一个死信表。 - 理由:Eventarc 提供了基于 CloudEvents 的事件驱动编排,可按存储桶和对象前缀进行过滤。批处理 Dataflow 作业在加载有效数据的同时,分离出格式错误的行进行分析。
- 保护平台安全
- 操作:使用不同的服务账号:生产者在
telemetry主题上获得pubsub.publisher权限;Dataflow worker SA 在telemetry-stream-sub上获得pubsub.subscriber权限,并获得对目标 BigQuery 数据集和 Cloud Storage 的写入权限;Eventarc 触发器使用一个专用的 SA,该 SA 在 Cloud Run 上具有invoker权限。在telemetry主题和 BigQuery 数据集上启用 CMEK。为所有推送端点(如有)配置 OIDC 和受众检查。 - 理由:最小权限 IAM 和 CMEK 满足了安全与合规要求;认证交付可防止欺骗。
- 可靠地运维和扩缩
- 操作:设置 Dataflow 自动扩缩,并配置一个足够大的
maxWorkers以吸收峰值负载。监控subscription/oldest_unacked_message_age和expired_ack_deadline_count;当超过阈值时发出警报。对于破坏兼容性的管道变更,使用drain模式进行部署以避免消息丢失。如果延迟增长,增加订阅者的并行度,并按处理时间的比例延长 ack 截止时间。 - 理由:主动监控能及早发现延迟和重传。自动扩缩和调优的 ack 截止时间可防止重复消息风暴。排空(Draining)操作可在升级期间保留处理中的消息。
该设计提供了一个弹性、安全且可观测的实时摄取方案,集成了事件驱动的批处理,支持重复容忍和 schema 演进,并能在隔离错误数据的同时提供快速分析以进行有针对性的修复。
← 使用 Dataflow 和 Apache Beam 进行流处理 · 所有领域 · Spark、Dataproc 与分布式数据处理 →
练习这些题目 → · 在 ExamRoll.io 上限时练习 →
Pass the whole exam — not just this question
You found this answer. Get every verified question and explanation in one place, and save hours of prep. Free to start.
通过考试 →