Skip to content
bitzorcas
中EN

Guide

Reporting 工单事件投影与乱序一致性

从 Ticket 聚合事件、生成式 wire payload 和 CAP consumer 出发,说明六类生命周期投影、最终 ID 缺陷、LastEventId 幂等、时间排序、重试、并发 Upsert 与恢复要求。

Last updated

TicketReportingEventConsumer 把 Tickets 生命周期事件折叠成 Rpt_TicketSummary 当前快照。topic、业务 payload 和失败重试链已经接通,六类状态也有 consumer;当前阻断点转移到了 opened 事件的实体 ID、缺少聚合版本,以及 Store 的并发写入语义。

1. 发布与消费契约

Tickets 聚合抛出的六个领域事件都继承 DomainEvent、实现 IIntegrationEvent 并声明 [IntegrationTopic]。源生成器建立事件类型到 topic 和业务标量字段的映射;DomainEventDispatchPipelineBehavior 再补上框架字段:

CAP payload 的组成
eventId = DomainEvent.EventId
eventType = 领域事件类型名
业务字段 = 源生成 mapper 从 Ticket* 事件读取的全部公开属性

该消息由事务内 INotificationPublisher 写入 CAP Outbox。发布异常或提交异常向外传播,业务数据与 Outbox 一起回滚。

Reporting consumer 使用独立的绑定 DTO。只要字段名和类型保持一致,CAP 可以把同一 wire body 绑定为下列类型:

TopicTickets 发布类型Reporting 参数类型业务字段
ticket.openedEvents.TicketOpenedTicketOpenedIntegrationEventTicketId、TenantId、RequesterId、Subject、PriorityName、OccurredAt
ticket.assignedEvents.TicketAssignedTicketAssignedIntegrationEventTicketId、TenantId、AssigneeId、AssignedBy、OccurredAt
ticket.startedEvents.TicketStartedTicketStartedIntegrationEventTicketId、TenantId、StartedBy、OccurredAt
ticket.resolvedEvents.TicketResolvedTicketResolvedIntegrationEventTicketId、TenantId、ResolvedBy、OccurredAt
ticket.closedEvents.TicketClosedTicketClosedIntegrationEventTicketId、TenantId、ClosedBy、OccurredAt
ticket.reopenedEvents.TicketReopenedTicketReopenedIntegrationEventTicketId、TenantId、ReopenedBy、OccurredAt
Rpt_TicketSummaryReporting consumerCAP OutboxGenerated topic registryTicket aggregateRpt_TicketSummaryReporting consumerCAP OutboxGenerated topic registryTicket aggregateTicket* DomainEvent + IIntegrationEventtopic + eventId/eventType + business fieldsbind Ticket*IntegrationEventread, order check, Upsert

当前测试直接调用 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 未使用。结果是:

  1. opened 消息把 TicketId="0" 写入 Mart;
  2. 保存完成后的 Ticket 响应和 assigned/started/… 事件使用最终雪花 ID;
  3. 后续 consumer 找不到对应 opened 行,抛出 “cannot be projected before its opened event”;
  4. 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 内部实现。恢复面可以选择受控来源:

  1. 版本化事件归档;
  2. Tickets owner 提供快照导出流;
  3. 受控 replica/ETL;
  4. shadow Mart 重建后做 count、状态分布和抽样对账,再原子切换。

当前仓库没有 Reporting inbox、checkpoint、事件归档适配器、重建命令或对账任务。“可重建”仍是目标能力。

10. 必测投影矩阵

  1. OpenTicket 保存后 opened 的 TicketId 是最终 ID,不是 0。
  2. 六个真实事务 Outbox payload 均可绑定 consumer DTO。
  3. A,A;A,B,A;相同时间戳;更早时间;version gap。
  4. non-opened 先到、opened 延迟到达和 opened 永久缺失。
  5. Mart read/upsert 临时失败、重试耗尽和 DLQ。
  6. 两 consumer 并发同 Ticket,first insert 唯一冲突。
  7. Reopened 明确清空 ClosedAt。
  8. TenantId/ TicketId 组合隔离和非法 payload。
  9. inbox 与 Mart 更新之间崩溃恢复。
  10. live projection 与 shadow rebuild 结果一致。

11. 审查命令

Terminal window
# 发布类型、生成式 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'

返回 Reporting 总览 · Tickets 事件与通知 · 测试与 GA

100%

滚轮或按钮缩放 · 放大后拖动画面 · 双击切换 100% / 200%