agents24 仓库 projection-patterns 技能全解:为事件溯源与 CQRS 系统构建投影与读模型
agents24 仓库 projection-patterns 技能全解为事件溯源与 CQRS 系统构建投影与读模型【免费下载链接】agentsMulti-harness agentic plugin marketplace for Claude Code, Codex, Cursor, OpenCode, GitHub Copilot, and Google Antigravity项目地址: https://gitcode.com/GitHub_Trending/agents24/agents本指南深度解读 agents24/agents 仓库中backend-development插件所附带的projection-patterns技能。该技能面向事件溯源 CQRS架构下的读侧实现从事件流构建读模型Read Model、物化视图Materialized View与搜索索引并把查询性能与实时报表作为目标场景。读完本文你将掌握投影器Projector的四种工作模式Live / Catchup / Persistent / Inline的选择依据、幂等投影与检查点恢复的工程要点并拿到 5 份可直接落地的 Python 模板订单汇总、Elasticsearch 搜索索引、日销售聚合、多表客户活动投影等以及对应的抽象基类与调度器实现。一、技能定位与使用场景projection-patterns是一个通过 YAML frontmatter 声明的 Agent Skill其description字段明确给出了触发条件Use when implementing CQRS read sides, building materialized views, or optimizing query performance in event-sourced systems即当你要为事件溯源系统实现 CQRS 读侧、构建物化视图或优化查询性能时应主动加载该技能。这与仓库的编写规范见 docs/authoring.md 中关于 Description triggers 的约定保持一致——描述中的触发短语正是模型决定是否调用该技能的依据。技能正文SKILL.md开篇即点明其主题构建投影Projection与读模型Read Model的全面指南。结合 frontmatter该技能适合以下六类工作构建 CQRS 读模型从事件创建物化视图优化查询性能实现实时仪表盘从事件构建搜索索引跨事件流聚合数据。从仓库的组织结构看该技能位于backend-development插件之下安装方式见 docs/plugins.md例如/plugin install backend-development与cqrs-implementation、event-store-design、saga-orchestration等技能共同构成一套完整的事件驱动后端技能族分别对应 event-sourcing-architect.md 中列出的能力项Projection building and read model optimization、CQRS (Command Query Responsibility Segregation) patterns、Event store design and implementation 等。也就是说投影是这个 Agent 在实现read side时最关键的一环其工作流第 4 步就是Build projections for query requirements。二、投影核心架构2.1 三段式数据管线技能用一张 ASCII 图刻画了投影的基本架构其链路为Event Store ──► Projector ──► Read Model (Database)Event Store事件存储存放不可变事件的追加式日志。每个事件是已经发生的事实只增不改不删且具备全局顺序即每个事件对应一个global_position。这一侧的设计细节如 Append-only、Ordered、Versioned、Subscriptions、Idempotent 等存储需求以及 EventStoreDB / PostgreSQL / Kafka / DynamoDB / Marten 等选型权衡详见配套技能 event-store-design/SKILL.md。Projector投影器读取事件流按事件类型分发给对应的 Handler 逻辑把状态变化应用到读模型。它相当于事件流与查询数据库之间的翻译器。Read Model读模型落库物化后的查询视图可能是表Tables、视图Views或缓存Cache。为查询而反范式化存储是 CQRS 读侧的核心产出。这三个角色在 cqrs-implementation/SKILL.md 的架构图中同样出现Write Model 通过 Events 单向流向 Read Model而Projector正是Updates read model from events的组件。投影技能与 CQRS 技能互为表里CQRS 决定读写分离的整体形状投影则回答事件到底如何变成可查询的数据。2.2 四种投影类型的选择类型说明适用场景Live实时通过订阅持续消费新事件当前状态查询、实时仪表盘Catchup追赶处理存量历史事件重建读模型、首次全量构建Persistent持久化存储检查点checkpoint可续跑重启后从断点恢复Inline内联与写操作处于同一事务需要强一致性要点在于Live 保证低延迟但不保证断点续传Persistent 用检查点换取可靠恢复二者常常配合使用Live 订阅的同时周期性持久化位置Catchup 与 Inline 是互补的两个极端——前者为了吞吐量牺牲一致性通常离线批量执行后者牺牲吞吐量换取写完即可读到适合一致性要求极高的少数场景。Catchup 通常也是重建持久化投影的第一步。整体上生产系统最常见的是Persistent Live的组合日常读查询都打向已就绪的物化读模型而不是事件存储本身。三、从抽象基类到投影调度器模板一精读技能正文仅 62 行是典型的导航层遵循仓库渐进式披露progressive disclosure约定见 docs/authoring.mdSKILL.md 只保留导航与快速开始完整模板与精讲案例下沉到references/details.md按需加载。因此模板库与详尽的实操示例都在 references/details.md共包含 5 个模板。下面逐一拆解并做源码级展开。3.1 Event 数据模型与 Projection 抽象基类第一份代码定义了两个基础抽象。先看事件的数据载体from abc import ABC, abstractmethod from dataclasses import dataclass from typing import Dict, Any, Callable, List import asyncpg dataclass class Event: stream_id: str event_type: str data: dict version: int global_position: int该Event结构完整刻画了事件溯源的核心字段stream_id所属事件流通常对应一个聚合实例配合 event-store-design 的建议可用Order-{uuid}这类带聚合类型前缀的 IDevent_type事件类型名投影器靠它做路由分发如OrderCreated、OrderShippeddata事件载荷payload应为与事件名对应的最小领域数据version聚合内版本号配合乐观并发控制global_position全局位置是投影检查点checkpoint的基础——所有投影都根据它知道自己已消费到哪里。接着是每个投影都必须实现的抽象基类class Projection(ABC): Base class for projections. property abstractmethod def name(self) - str: Unique projection name for checkpointing. pass abstractmethod def handles(self) - List[str]: List of event types this projection handles. pass abstractmethod async def apply(self, event: Event) - None: Apply event to the read model. pass三个抽象成员分别回答三个问题name提供唯一投影名作为检查点存储中的键必须全局唯一模板里分别叫order_summary、product_search、daily_sales、customer_activityhandles()声明该投影关心的事件类型白名单投影器据此决定是否需要调用apply也可当作这个读模型由哪些事件驱动的自文档apply(event)是单事件处理器负责把一件事件落到读模型上。3.2 Projector 调度器检查点续跑与全量重建class Projector: Runs projections from event store. def __init__(self, event_store, checkpoint_store): self.event_store event_store self.checkpoint_store checkpoint_store self.projections: List[Projection] [] def register(self, projection: Projection): self.projections.append(projection) async def run(self, batch_size: int 100): Run all projections continuously. while True: for projection in self.projections: await self._run_projection(projection, batch_size) await asyncio.sleep(0.1) async def _run_projection(self, projection: Projection, batch_size: int): checkpoint await self.checkpoint_store.get(projection.name) position checkpoint or 0 events await self.event_store.read_all(position, batch_size) for event in events: if event.event_type in projection.handles(): await projection.apply(event) await self.checkpoint_store.save( projection.name, event.global_position ) async def rebuild(self, projection: Projection): Rebuild a projection from scratch. await self.checkpoint_store.delete(projection.name) # Optionally clear read model tables await self._run_projection(projection, batch_size1000)这个调度器虽短却浓缩了投影运行时的全部关键机制值得逐段展开1. 注册制与轮询式调度。register()把多个投影挂到同一 Projector 上run(batch_size100)进入while True循环逐个投影调用_run_projection随后asyncio.sleep(0.1)节流避免对事件存储造成忙轮询。0.1 秒的间隔就是该实现的近实时下限在真正生产环境里这里通常会被事件存储的推送订阅subscription替代驱动方式从 pull 变为 push。2. 检查点是可靠性的核心。_run_projection首先从checkpoint_store.get(projection.name)读取该投影上次处理到的位置取不到则从 0事件流开头开始。这正是四种类型中Persistent的实现形态每处理一个事件就用event.global_position更新检查点。因此进程崩溃、重启或部署发布后投影能从断点继续而不是从零重放。注意实现细节上的取舍代码是处理一个、保存一个若追求更高吞吐可改为每批结束时统一保存一次检查点代价是崩溃时最多回退一个批次、可能重复处理若干事件——这又反过来要求apply是幂等的见下文最佳实践。3. 内建类型过滤。分发逻辑不是让每个投影收到全部事件再去判断而是if event.event_type in projection.handles()事件类型不在白名单就跳过但检查点仍然前移保证不同类型的事件不会阻塞某个投影的位置推进。4.rebuild是重建能力的模板。删除检查点注释提示还需要按需清空读模型表然后用更大的批次1000从零重新跑一遍。这与技能 Best Practices 中的Plan for rebuilds直接呼应——因为事件日志不可变读模型在理论上永远可以从事件流重新长出来这正是事件溯源架构让读侧具有可重构性的根因。从工程语义上看run()对应Live/Catchup的统一形态传入batch_size决定单批拉取量连续消费属于 Live而rebuild()以更大批次从 0 消费、本身就是一个显式的 Catchup 操作。四、五个模板的实战拆解4.1 模板二Order Summary 投影单行聚合式读模型OrderSummaryProjection监听订单生命周期中的六个事件OrderCreated、OrderItemAdded、OrderItemRemoved、OrderShipped、OrderCompleted、OrderCancelled把订单的实时状态折叠进order_summaries表的单一行。apply()内部用一张事件类型 → 处理方法的字典做二次路由比 if/else 链更易扩展async def apply(self, event: Event) - None: handlers { OrderCreated: self._handle_created, OrderItemAdded: self._handle_item_added, OrderItemRemoved: self._handle_item_removed, OrderShipped: self._handle_shipped, OrderCompleted: self._handle_completed, OrderCancelled: self._handle_cancelled, } handler handlers.get(event.event_type) if handler: await handler(event)各处理器的语义可概括为一张状态转移—SQL 副作用对照表事件读模型副作用OrderCreated插入一行order_summaries初始statuspending、total_amount0、item_count0OrderItemAdded行内增量更新total_amount price*quantityitem_count 1OrderItemRemoved反向扣减total_amount - price*quantityitem_count - 1OrderShipped置statusshipped并写入shipped_atOrderCompleted置statuscompleted并写入completed_atOrderCancelled置statuscancelled记录cancelled_at与cancellation_reason用event.data.get(reason)容错从这段代码可以提炼出三个值得学习的设计选择把事件语义翻译成增量更新读模型不需要知道订单的完整历史只需要在每个事件到来时做一次局部更新如加价、减数、改状态。事件驱动投影的本质就是把INSERT/UPDATE/DELETE的副作用散进各个事件处理器里读查询因此永远 O(1) 命中目标行。状态机字段各自独立推进status、shipped_at、completed_at、cancelled_at分别由各自事件更新updated_at NOW()统一维护。若担心事件乱序如OrderCancelled早于OrderCompleted到达实际生产还需要在更新前比较事件时间戳或版本——这正是技能Dont ignore ordering告诫的来源。异步非阻塞写通过asyncpg.Pool连接池acquire()执行 SQL与整个技能模板的 asyncio 风格保持一致支撑高并发投影。4.2 模板三Elasticsearch 搜索投影读模型即搜索引擎索引ProductSearchProjection演示了读模型不一定是 SQL 表——它把产品事件投影为 Elasticsearch 的products索引让查询侧获得全文检索能力正是技能 When-to-use 中Building search indexes from events的落地形态async def apply(self, event: Event) - None: if event.event_type ProductCreated: await self.es.index( indexself.index, idevent.data[product_id], document{ name: event.data[name], description: event.data[description], category: event.data[category], price: event.data[price], tags: event.data.get(tags, []), created_at: event.data[created_at] } ) elif event.event_type ProductUpdated: await self.es.update( indexself.index, idevent.data[product_id], doc{ name: ..., description: ..., updated_at: ... } ) elif event.event_type ProductPriceChanged: await self.es.update( indexself.index, idevent.data[product_id], doc{ price: event.data[new_price], price_updated_at: event.data[changed_at] } ) elif event.event_type ProductDeleted: await self.es.delete(indexself.index, idevent.data[product_id])四个事件分别映射为索引层面的 create / update / partial-update / delete且ProductPriceChanged只做字段级局部更新doc里只有price避免了全文档重写。两个细节尤其值得注意document只装载查询需要的字段名称、描述、分类、价格、标签、时间是对反范式化、为查询优化的贯彻——事件里的原始字段可能很多读模型只需保留检索与展示所需子集用event.data.get(tags, [])提供默认值体现对缺失字段的防御式处理。4.3 模板四跨流聚合投影Daily Sales 报表DailySalesProjection演示了聚合数据这一类投影它把两条完全不同的事件流OrderCompleted与OrderRefunded汇聚到按日期分组的daily_sales表async def _increment_sales(self, event: Event): date event.data[completed_at][:10] # YYYY-MM-DD async with self.pool.acquire() as conn: await conn.execute( INSERT INTO daily_sales (date, total_orders, total_revenue, total_items) VALUES ($1, 1, $2, $3) ON CONFLICT (date) DO UPDATE SET total_orders daily_sales.total_orders 1, total_revenue daily_sales.total_revenue $2, total_items daily_sales.total_items $3, updated_at NOW() , date, event.data[total_amount], event.data[item_count] )这里的核心技巧是 PostgreSQL 的INSERT ... ON CONFLICT (date) DO UPDATEUPSERT某天的第一条完成订单触发 INSERT后续订单全部走冲突分支做增量累加从而在单条语句内同时处理该天首单与该天后续订单两种情况天然幂等且免去先查后写。total_orders 1、total_revenue $2、total_items $3的写法全部是增量式refunded侧则反向递减订单数与收入同时单独累计total_refunds便于对账。4.4 模板五多表事务投影Customer ActivityCustomerActivityProjection展示了单事件触发跨多张表写入的场景并显式启用事务async def apply(self, event: Event) - None: async with self.pool.acquire() as conn: async with conn.transaction(): if event.event_type CustomerCreated: # Insert into customers table await conn.execute( INSERT INTO customers (customer_id, email, name, tier, created_at) VALUES ($1, $2, $3, bronze, $4), ...) # Initialize activity summary await conn.execute( INSERT INTO customer_activity_summary (customer_id, total_orders, total_spent, total_reviews) VALUES ($1, 0, 0, 0), ...) elif event.event_type OrderCompleted: ...该投影监听CustomerCreated、OrderCompleted、ReviewSubmitted、CustomerTierChanged四个事件分别维护三类数据customers主表创建客户、更新会员等级tier新用户默认bronzecustomer_activity_summary汇总表初始化全 0随后累加订单数、消费额、评价数并记录last_order_at/last_review_atcustomer_order_history明细流水表追加每笔完成订单。这正是技能 Best Practices 中Use transactions for multi-table updates的直接体现一次apply内若涉及两张表的写操作就放进asyncpg的conn.transaction()保证要么全部提交、要么全部回滚读模型内部永远一致。同理CustomerCreated同时写主表与汇总表如果不用事务就有可能出现客户已建、汇总行缺失的中间状态。五、最佳实践Dos 与 Donts技能把经验浓缩为五条 Dos 与四条 Donts结合上文模板可一一对应验证Dos应该做Make projections idempotent投影必须幂等可安全重放事件重放replay是重建、故障恢复的常态。幂等的直接收益是同一条事件处理两遍也不会破坏数据。上文模板中的 UPSERT、增量更新天然幂等重复执行同一事件相当于再多加一次注意需要配合事件去重或检查点机制凡是做INSERT的地方都要问一句重放时会不会主键冲突必要时改成 UPSERT。Use transactions多表更新使用事务见模板五async with conn.transaction()保证跨表原子性。Store checkpoints持久化检查点见模板一调度器每事件或每批把global_position写入 checkpoint 存储是重启后无缝续跑的前提。Monitor lag监控投影延迟投影落后于事件流意味着读模型数据过期。生产上应为每个投影记录已消费位置 vs 事件流最新位置的差值并设置告警——这与仓库中 observability-monitoring 等运维插件的思路一致滞后本身是系统健康状况的信号。Plan for rebuilds为重建做设计读模型必须可丢弃、可重建。rebuild()展示了删检查点、清表、重放三步走正因为事件日志是不可变的系统事实对应 event-sourcing-architect.md 中Events are facts - never delete or modify them重建才永远可行。Donts不要做Dont couple projections投影之间不要耦合每个投影独立消费事件、独立维护检查点、独立失败与重建。一个投影的故障不应拖垮其他投影。Projector 中每个投影独立_run_projection即此原则的实现。Dont skip error handling不要省略错误处理消费失败要记录日志并告警而不是静默跳过——否则检查点继续前移会造成读模型永久缺数据。Dont ignore ordering不要忽视事件顺序事件必须按序处理。global_position天然提供全序若并行消费同一事件流必须保证每个投影内的顺序或通过乐观锁/版本比较拒绝过期事件。Dont over-normalize不要过度范式化读模型要按查询模式反范式化denormalize把 join、聚合提前到写入时完成查询侧才能简单快速。order_summaries单行冗余了总额与条目数、daily_sales预聚合了日报都是这个原则的例子。六、如何在真实项目中落地这套技能6.1 技能在仓库中的装配方式projection-patterns是 Agent Skill 形式的分层文档导航层 SKILL.md 保持轻量约 60 行低于仓库对 Codex 8 KB 体积上限的约束深度模板在按需加载的 references/details.md。当 Agent 判定任务属于 CQRS 读侧 / 物化视图 / 查询优化时先读 SKILL.md 建立上下文再按需展开 references/details.md 拿模板这正是仓库progressive disclosure的渐进披露约定docs/authoring.md。它通常与cqrs-implementation读写作分离的总体设计、event-store-design事件存储选型与建模配合使用而backend-development插件的feature-development编排命令feature-development.md在实现事件驱动类特性时会先由backend-development-backend-architect等 Agent 产出架构设计其中的读侧实现即依赖此类技能。也就是说若你在 Claude Code、Codex、Cursor、OpenCode 或 Antigravity 等 harness 中安装了 backend-development 插件并规划事件溯源系统这套技能会在实现读模型时自动被引用。6.2 把模板接入真实系统的建议顺序选型事件存储确认你的事件流能提供全局单调递增的位置PostgreSQL 序列、EventStoreDB 的 global position、Kafka 的分区偏移等均可这是检查点机制的前提。按查询需求设计读模型列出每个查询场景需要的字段与聚合粒度反范式化建模订单汇总单行、日销售按日期分组、客户活动多表冗余据此选择要监听的event_type集合。实现 Projection 子类继承Projection实现name/handles()/apply()把事件→SQL/索引副作用写进各事件处理器遵循状态增量更新与幂等原则。实现 checkpoint 存储模板中checkpoint_store为抽象接口生产可用独立表projection_name PRIMARY KEY, position BIGINT、Redis 或对象存储关键是一次更新要快且可靠。接线 Projectorregister()注册所有投影后启动run()写一个离线脚本或运维命令按需调用rebuild()并配置 lag 监控告警。七、小结投影的本质是把不可变事件流翻译成可查询、可聚合、可检索的读模型从而让事件溯源系统的查询侧获得与业务写模型解耦的性能与灵活性。projection-patterns技能用一张架构图、一张类型对照表、一套抽象基类与调度器、五份 Python 模板和九条 Do/Dont 经验完整覆盖了从理解投影架构到写出可运行投影器再到上线后可靠运维的全过程。把它与同族的 cqrs-implementation读写分离和 event-store-design事件存储技能配合使用即可在事件驱动后端中搭建出一套具备检查点恢复、幂等重放、事务化多表写入与全量重建能力的完整读侧体系。【免费下载链接】agentsMulti-harness agentic plugin marketplace for Claude Code, Codex, Cursor, OpenCode, GitHub Copilot, and Google Antigravity项目地址: https://gitcode.com/GitHub_Trending/agents24/agents创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
