TicketReportingEventConsumer 把 Tickets 生命周期事件折叠成 Rpt_TicketSummary 当前快照。topic、业务 payload 和失败重试链已经接通,六类状态也有 consumer;当前阻断点转移到了 opened 事件的实体 ID、缺少聚合版本,以及 Store 的并发写入语义。
1. 发布与消费契约
Tickets 聚合抛出的六个领域事件都继承 DomainEvent、实现 IIntegrationEvent 并声明 [IntegrationTopic]。源生成器建立事件类型到 topic 和业务标量字段的映射;DomainEventDispatchPipelineBehavior 再补上框架字段:
eventId = DomainEvent.EventIdeventType = 领域事件类型名业务字段 = 源生成 mapper 从 Ticket* 事件读取的全部公开属性该消息由事务内 INotificationPublisher 写入 CAP Outbox。发布异常或提交异常向外传播,业务数据与 Outbox 一起回滚。
Reporting consumer 使用独立的绑定 DTO。只要字段名和类型保持一致,CAP 可以把同一 wire body 绑定为下列类型:
| Topic | Tickets 发布类型 | Reporting 参数类型 | 业务字段 |
|---|---|---|---|
ticket.opened | Events.TicketOpened | TicketOpenedIntegrationEvent | TicketId、TenantId、RequesterId、Subject、PriorityName、OccurredAt |
ticket.assigned | Events.TicketAssigned | TicketAssignedIntegrationEvent | TicketId、TenantId、AssigneeId、AssignedBy、OccurredAt |
ticket.started | Events.TicketStarted | TicketStartedIntegrationEvent | TicketId、TenantId、StartedBy、OccurredAt |
ticket.resolved | Events.TicketResolved | TicketResolvedIntegrationEvent | TicketId、TenantId、ResolvedBy、OccurredAt |
ticket.closed | Events.TicketClosed | TicketClosedIntegrationEvent | TicketId、TenantId、ClosedBy、OccurredAt |
ticket.reopened | Events.TicketReopened | TicketReopenedIntegrationEvent | TicketId、TenantId、ReopenedBy、OccurredAt |
当前测试直接调用 consumer DTO,尚未通过真实事务 Outbox 捕获字节再执行 CAP model binding。因此“字段在源码层对齐”已有证据,“宿主序列化后的字节契约”仍需独立集成测试。
2. 六类状态变化
| 事件 | 必须已有 opened 行 | 写入结果 |
|---|---|---|
| opened | 否 | 新建 New;写 RequesterId、Subject、PriorityName、OpenedAt |
| assigned | 是 | 写 AssigneeId,状态改为 Assigned |
| started | 是 | 状态改为 InProgress |
| resolved | 是 | 状态改为 Resolved |
| closed | 是 | 状态改为 Closed,ClosedAt=OccurredAt |
| reopened | 是 | 状态改为 Reopened,ClosedAt=null |
AssigneeId 现在有真实 writer;Reopened 也能清空 ClosedAt,因为 Store 已改成直接 existing.ClosedAt = summary.ClosedAt。旧文档所述“null 保留旧关闭时间”和“只投影 opened/closed”均已过时。
3. Opened 事件仍使用占位 ID
Ticket.Open 的当前顺序是:
var ticket = new Ticket("0", tenantId, requesterId, subject, priority, nowUtc);
ticket.Raise(new TicketOpened( ticket.Id, // 此处仍为 "0" ticket.TenantId, ticket.RequesterId, ticket.Subject, ticket.Priority.Name, nowUtc));
// OpenTicket Handler 随后调用 TicketRepository.SaveAsync(ticket)。// IEntitySet<T>.AddAsync 才会通过 persistence adapter 分配最终 ID。实体基类已经提供 RaiseWhenPersistenceIdAssigned(Func<TId, IDomainEvent>),但 Ticket.Open 未使用。结果是:
- opened 消息把
TicketId="0"写入 Mart; - 保存完成后的 Ticket 响应和 assigned/started/… 事件使用最终雪花 ID;
- 后续 consumer 找不到对应 opened 行,抛出 “cannot be projected before its opened event”;
- CAP 重试不会自动修复 ID 不一致。
应把 opened 事件改为 RaiseWhenPersistenceIdAssigned(id => new TicketOpened(id, ...)),并增加从 OpenTicketCommand、持久化、Outbox 到 Reporting 的端到端测试。不要在 consumer 里把 0 猜测成某个最新工单。
4. 重复、陈旧和前置缺口
非 opened 路径的共同决策如下:
var existing = await GetExistingAsync(tenantId, ticketId, eventName, ct);
// LastEventId 只识别当前最后一次写入的直接重复。if (existing.LastEventId == eventId) return; // 直接重复
if (occurredAt < existing.LastUpdatedAt) return; // 更早时间的事件
// 相同时间戳不会被视为 stale;继续覆盖。// 因此生产顺序仍需要 Version/sequence 提供确定性。Apply(existing);existing.LastEventId = eventId;existing.LastUpdatedAt = occurredAt;await UpsertAsync(existing, eventName, ct);opened 会检查直接重复,但不会比较 OccurredAt;它会构造完整 New snapshot。非 opened 在缺少前置行时抛异常,让 CAP 重试,而不是确认并丢弃。这个变化修复了旧版“Closed 先到后永久丢失”的故障模式,但没有 gap ledger;如果 opened 永远无法成功,消息最终仍会进入重试耗尽状态。
5. LastEventId 与 OccurredAt 的边界
LastEventId 只识别最后事件的直接重投:A,A。序列 A,B,A 中,第二个 A 不再等于 LastEventId;是否跳过完全取决于 OccurredAt。若重放的 A 时间更早,会被时间门禁跳过;若时间相同或错误地更新,则仍可能覆盖。
OccurredAt 也不是严格业务序号:
- 两个状态变化可能拥有相同时间精度;当前比较使用
<而不是<=; - 服务时钟、导入数据或手工修复可能产生相同或异常时间;
- 它不能识别 version gap;
- opened 本身不执行时间门禁。
生产级投影应在事件中携带每个 Ticket 单调递增的 AggregateVersion,并在 Mart 保存 LastAppliedVersion。规则应为:
version <= current:duplicate/stale,确认并计指标;version == current + 1:原子应用;version > current + 1:记录 gap,等待回补,不静默越过。
EventId inbox 与版本门禁解决不同问题:inbox 防任意重投,version 决定业务顺序。两者的写入应与 Mart mutation 同事务。
6. 失败与 CAP 重试
consumer 的 catch 会记录异常类型并重新抛出;以下情况不会伪装成消费成功:
- Mart read Result 失败;
- Mart Upsert Result 失败;
- non-opened 找不到前置行;
- 其他非取消、非内存不足异常。
OperationCanceledException 和 OutOfMemoryException 不在 catch filter 内,也会自然传播。日志不输出 Subject 或完整 payload。
仍应补充明确错误分类:
| 类别 | 建议处理 |
|---|---|
| 临时数据库/网络失败 | 有界重试,超限后 DLQ/告警 |
| 缺少 opened 前置 | 记录 gap;短期重试;可受控回源 |
| schema/必填字段非法 | quarantine,不进行无界重试 |
| duplicate/stale | 确认成功并计数 |
| 并发版本冲突 | 重新读取并按版本重试 |
7. Store 并发窗口
ReportingMartStore.UpsertTicketSummaryAsync 仍是:先按 (TenantId,TicketId) 查询,再 Add 或 Update。它没有原子 Upsert、ExpectedVersion 或 affected-row compare-and-set。
因此可能出现:
- 两个 first event 都查不到,随后在唯一键上竞争;
- 两个 consumer 读取同一旧行,后提交的写入覆盖先提交的更新;
- LastEventId、LastUpdatedAt 与业务字段之间没有数据库条件;
- 同时间戳事件以提交顺序决定最终状态。
数据库唯一索引防止重复物理行,但不能保证正确业务顺序。双 ORM 测试应覆盖唯一冲突恢复和条件更新,不应以 InMemory Store 代替。
8. 当前测试证据
TicketReportingEventConsumerTests 已验证:
- Opened→Assigned→Started→Resolved→Closed→Reopened;
- AssigneeId 和六个状态名称;
- ClosedAt 写入与 Reopened 清空;
- 相邻 Reopened 重复不再次 Upsert;
- 较早 Closed 不倒退 Reopened 快照。
尚未验证:
Ticket.Open最终 ID;- 源生成 payload 的真实 CAP 字节绑定;
- A,B,A、相同 OccurredAt、跨版本 gap;
- 缺前置消息的重试/DLQ 策略;
- 两 consumer 并发和真实数据库唯一冲突;
- SqlSugar/EF Core 等价行为;
- replay、reconcile 与 shadow rebuild。
9. 重建与回源
普通 consumer 不应反向调用 Tickets 内部实现。恢复面可以选择受控来源:
- 版本化事件归档;
- Tickets owner 提供快照导出流;
- 受控 replica/ETL;
- shadow Mart 重建后做 count、状态分布和抽样对账,再原子切换。
当前仓库没有 Reporting inbox、checkpoint、事件归档适配器、重建命令或对账任务。“可重建”仍是目标能力。
10. 必测投影矩阵
- OpenTicket 保存后 opened 的 TicketId 是最终 ID,不是
0。 - 六个真实事务 Outbox payload 均可绑定 consumer DTO。
- A,A;A,B,A;相同时间戳;更早时间;version gap。
- non-opened 先到、opened 延迟到达和 opened 永久缺失。
- Mart read/upsert 临时失败、重试耗尽和 DLQ。
- 两 consumer 并发同 Ticket,first insert 唯一冲突。
- Reopened 明确清空 ClosedAt。
- TenantId/ TicketId 组合隔离和非法 payload。
- inbox 与 Mart 更新之间崩溃恢复。
- live projection 与 shadow rebuild 结果一致。
11. 审查命令
# 发布类型、生成式 payload 和六个订阅入口。rg -n "IntegrationTopic|CapSubscribe|BuildParameters|Ticket(Open|Assign|Start|Resolv|Clos|Reopen)" \ src/Platform/{Tickets,Reporting} src/Framework/BitzOrcas.Application/Pipelines -g '*.cs'
# 最终 ID 缺陷必须在修复后由 RaiseWhenPersistenceIdAssigned 取代直接 Raise。rg -n 'new Ticket\(|Raise\(new TicketOpened|RaiseWhenPersistenceIdAssigned' \ src/Platform/Tickets -g '*.cs'
# 幂等、时间门禁、并发 Upsert 与测试证据。rg -n "LastEventId|LastUpdatedAt|FirstOrDefaultAsync|TicketReportingEventConsumerTests" \ src/Platform/Reporting tests -g '*.cs'