site logo

Marico's space

GCP事件驱动架构实践:Pub/Sub、Cloud Tasks与Cloud Scheduler

服务器技术 2026-08-26 14:49:12 10

最近在搞AI评估平台,踩了不少坑,终于把事件驱动这套东西跑通了。这篇把GCP上Pub/Sub、Cloud Tasks和Cloud Scheduler的实际用法说清楚,不整虚的。

为什么需要事件驱动

AI评估平台有个问题,用REST API根本没法优雅地解决:一个评估任务跑5到30分钟,涉及6个不同服务,中间随时可能挂。

用同步HTTP:超时、重试撞上已经完成的任务、调试噩梦。换成事件驱动:每个服务独立订阅。一个挂了,从自己的检查点重试就行。其他的根本不知道出过问题。

管道阶段设计

六个阶段,每个阶段代表一次交接,下一个服务可能独立失败。

  • 阶段A:创建任务
  • 阶段B:数据预处理
  • 阶段C:执行Agent
  • 阶段D:收集输出
  • 阶段E:评判打分
  • 阶段F:汇总学习

每个阶段发布到Pub/Sub主题,下一个阶段订阅。如果阶段C挂了,阶段D根本不知道有任务被尝试过。阶段C从检查点重试,成功后阶段D再接着处理。这种独立性才是事件驱动架构的精髓。

Phase A: Task Creation → (Pub/Sub: task-events)
Phase B: Preprocessing → (Pub/Sub: phase-b-batch-notifications)
Phase C: Agent Execution → (Pub/Sub: phase-c-asset-feed)
Phase D: Output Collection → (Pub/Sub: phase-e-jobs)
Phase E: Judging & Scoring → (Pub/Sub: phase-f-jobs)
Phase F: Aggregation & Learning → (Pub/Sub: evaluation-events)
Consumers: Notification, RAG Sync, Analytics

扇出模式

这是最强大的模式:一个主题,多个订阅。evaluation-events只发布一次,但被通知服务、RAG同步、分析管道和自动评估调度器同时消费。想加新消费者?建个订阅就行。发布者代码零改动。

# Publisher (evaluation engine)
publisher = pubsub_v1.PublisherClient()
topic_path = publisher.topic_path(project_id, "evaluation-events")
publisher.publish(topic_path, json.dumps({ "evaluation_id": "eval_123", "status": "completed", "score": 0.85
}).encode()) # Multiple subscribers (each independent)
# - notification-service subscribes to send emails
# - rag-service subscribes to update embeddings
# - analytics-service subscribes to update dashboards
# - auto-eval-service subscribes to trigger next batch

每个订阅者独立处理消息。通知服务挂了,RAG服务照样处理。这种解耦才是扇出模式真正厉害的地方。

死信队列

每个生产环境主题都配了DLQ,重试策略是5次。每周还要做一次DLQ review。死信队列专门捕获那些非临时性的失败:

  • Schema迁移问题:消息格式变了,但队列里还有旧格式的消息
  • 配额限制:Firestore拒绝写入,因为日限额用完了
  • Cloud SQL连接故障:如果是数据库真挂了,临时失败会变成永久失败

没有DLQ,这些都会变成静默数据丢失。消息失败、重试、再失败、然后消失。有了DLQ,每个失败都能看到,调查根因,修掉。

Cloud Scheduler触发自动化任务

每个定时任务都是往Pub/Sub主题发一条消息。消费服务根本不知道这消息是定时触发的、用户操作触发的、还是API调用触发的。想在生产环境暂停auto-eval-daily?不用改任何服务代码,直接停掉调度器就行。

# Cloud Scheduler job
gcloud scheduler jobs create pubsub auto-eval-daily \ --schedule="0 2 * * *" \ --topic=auto-eval-trigger \ --message-body='{"trigger": "daily", "batch_size": 100}' # Service handler (doesn't know the source)
def handle_auto_eval(message): data = json.loads(message.data) # Process the batch - same code path whether
 # triggered by schedule, API, or manual publish
 process_evaluation_batch(data["batch_size"])

这种抽象的好处:可以用手动发布来测试定时任务、停掉调度器来暂停任务、改调度时间不用部署代码。

监控

每个Cloud Run服务都配了错误率超过5%的告警。还监控CPU利用率低于10%的情况,用来自动检测闲置VM——这能省不少计算费用。但真正有价值的是告警策略:按错误率告警,不按单个错误告警。一条消息失败不是问题。错误率持续超过5%说明有东西坏了。还要监控DLQ里的消息积压时间——如果消息在DLQ里躺超过一小时,立刻调查。这种主动监控能在问题变事故之前就抓住它。

Pub/Sub还是Cloud Tasks:什么时候用哪个

Pub/Sub适合事件——发生的事情,其他服务可能关心。评估完成了。任务创建了。用户注册了。多个服务可能订阅同一个事件。

Cloud Tasks适合命令——要做的事情,要求Exactly-Once,有重试和截止时间。发一封邮件。处理一笔支付。更新一条数据库记录。一个任务,一个处理器,保证送达。

我们的评估管道用Pub/Sub,因为多个服务都关心评估事件。通知发送用Cloud Tasks,因为每封邮件都要保证只发一次,邮件服务down了就重试。区别很关键:Pub/Sub是发完就忘的事件。Cloud Tasks是可靠命令执行。

关键经验

消息顺序不重要——前提是你设计对了。每条消息都包含足够的上下文,可以独立处理。如果真需要顺序,在消息里带个序号或时间戳。别依赖Pub/Sub的顺序保证——那是尽力而为的,高负载下会断。

先提交再确认,别反了。先写数据库,再确认消息,永远是这样。如果确认完了再写数据库,写的时候挂了,消息直接没了。如果先写再确认,确认失败了,消息会重新投递——但幂等性检查会防止重复处理。

Schema演进不可避免。给消息Schema加版本号。每条消息都带version字段。想改Schema?同时发布新旧两个版本,逐步升级消费者,然后废弃旧版本。别用破坏性变更搞崩现有消费者。

Pub/Sub管事件,Cloud Tasks管命令。这个界限要清晰。事件是"发生了什么"。命令是"做这件事"。混着用会导致重试语义、送达保证、错误处理一团糟。