Java微服務事件與存儲設計說明(含代碼示例)
本文檔說明事件的結構、事件存儲的結構,以及事件的發(fā)送、保存、重試、關聯(lián)全流程。與 需求文檔 配套,作為實現(xiàn) bus / man 的數(shù)據(jù)與流程依據(jù)。
1. 事件的結構
1.1 概念區(qū)分
業(yè)務載荷(Payload):業(yè)務方關心的具體數(shù)據(jù),如訂單號、用戶 ID、操作類型等,結構由業(yè)務定義,系統(tǒng)以不透明字節(jié)或 JSON 存儲。
事件元數(shù)據(jù)(Metadata):系統(tǒng)為調度、追溯、重試而附加的字段,由 bus/man 統(tǒng)一約定。
下面所說的「事件」指系統(tǒng)層面的事件 envelope:元數(shù)據(jù) + 載荷。
1.2 事件統(tǒng)一結構(Envelope)
在 RabbitMQ 消息體與 MongoDB 持久化中,事件采用統(tǒng)一信封結構,便于發(fā)送、存儲與關聯(lián)。
| 字段 | 類型 | 必填 | 說明 |
|---|---|---|---|
| eventId | String | 是 | 全局唯一事件 ID(如 UUID),用于去重、重試、關聯(lián)。 |
| traceId | String | 推薦 | 分布式鏈路 ID,同一次業(yè)務觸發(fā)的多條事件/調用共享,用于串聯(lián)「事件調用鏈路」。 |
| spanId | String | 可選 | 當前環(huán)節(jié)的 span ID,與 traceId 一起構成鏈路節(jié)點。 |
| parentEventId | String | 可選 | 若本事件由另一事件的消費所觸發(fā),則填上游事件的 eventId,用于父子關聯(lián)。 |
| topic | String | 是 | 邏輯主題,對應 RabbitMQ 的 routing key 或交換機路由,如 order.created、payment.notify。 |
| payload | Object / String / Binary | 是 | 業(yè)務載荷,JSON 對象或序列化字符串;消費端按 topic 約定反序列化。 |
| payloadType | String | 可選 | 載荷類型提示,如 application/json、application/cloudevents+json。 |
| initiator | Object | 推薦 | 發(fā)起方信息,見下。 |
| occurredAt | String (ISO-8601) | 是 | 事件發(fā)生時間(業(yè)務發(fā)生時間或首次發(fā)送時間)。 |
| sentAt | String (ISO-8601) | 可選 | 首次發(fā)送到 MQ 的時間,由 bus 在發(fā)送時寫入。 |
| expireAt | String (ISO-8601) | 可選 | 業(yè)務過期時間,超過后可不再重試或不再展示。 |
initiator(發(fā)起方) 建議結構:
| 字段 | 類型 | 說明 |
|---|---|---|
| service | String | 發(fā)起服務名,如 order-service、bus。 |
| operation | String | 操作/接口,如 createOrder、POST /api/orders。 |
| userId | String | 若可識別,填用戶 ID。 |
| clientRequestId | String | 客戶端請求 ID,便于與業(yè)務日志關聯(lián)。 |
示例(JSON):
{
?"eventId": "evt-a1b2c3d4-e5f6-7890-abcd-ef1234567890",
?"traceId": "trace-xxx-001",
?"spanId": "span-001",
?"parentEventId": null,
?"topic": "order.created",
?"payload": { "orderId": "ORD-2024-001", "amount": 99.00 },
?"payloadType": "application/json",
?"initiator": {
? ?"service": "order-service",
? ?"operation": "createOrder",
? ?"userId": "user-123",
? ?"clientRequestId": "req-abc-001"
},
?"occurredAt": "2024-02-28T10:00:00Z",
?"sentAt": "2024-02-28T10:00:01Z"
}2. 事件存儲的結構(MongoDB)
2.1 集合設計
事件及鏈路數(shù)據(jù)存放在 同一邏輯庫(如 wx-bus),bus 與 man 均可訪問(或由 man 只讀 + 重推寫)。集合劃分如下:
| 集合名 | 用途 |
|---|---|
| events | 事件主表:每條事件一條文檔,僅包含 envelope + 發(fā)送/重試相關狀態(tài)(如 PENDING/SENT/RETRYING)。不在此文檔內嵌消費反饋,避免多消費者并發(fā)更新同一文檔導致寫沖突與狀態(tài)不一致。 |
| topic_consumers | Topic–消費者配置表:維護每個 topic 對應哪些消費者,支持動態(tài)配置(增刪改消費者無需發(fā)版)。事件創(chuàng)建時根據(jù)該 topic 的配置初始化 event_consumptions 的「預期消費者」記錄,見 2.3。 |
| event_consumptions | 消費反饋表:無 (eventId, consumerId) 唯一約束;同一消費者多次回調各保留一條記錄(attemptNo 遞增),便于問題排查。記錄來源:① 事件創(chuàng)建時按 topic_consumers 預插(attemptNo=0、success=null);② 消費者每次回調插入新記錄。通過 eventId 與 events 關聯(lián);匯總時按 (eventId, consumerId) 取最新一條參與計算,見 2.3、2.5、2.6。 |
| event_links(可選) | 若需將「鏈路」單獨建模為邊表,可存 (eventId, parentEventId, traceId, sequence);多數(shù)場景僅用 events 內字段即可。 |
這樣設計可避免:多個消費者同時回調時對同一 events 文檔做 read-modify-write 帶來的并發(fā)更新與狀態(tài)覆蓋問題。
2.2 events 文檔結構
在 1.2 的 envelope 基礎上,增加存儲與狀態(tài)字段:
| 字段 | 類型 | 說明 |
|---|---|---|
| _id | ObjectId | MongoDB 主鍵;也可用 eventId 作為 _id 便于冪等寫入。 |
| eventId | String | 同 1.2,唯一,建議建唯一索引。 |
| traceId | String | 同 1.2,建索引便于按鏈路查詢。 |
| spanId | String | 同 1.2。 |
| parentEventId | String | 同 1.2,建索引便于查「某事件的子事件」。 |
| topic | String | 同 1.2,必建索引(列表/篩選)。 |
| payload | Object / String | 同 1.2。 |
| payloadType | String | 同 1.2。 |
| initiator | Object | 同 1.2,可按 initiator.service 等建復合索引。 |
| occurredAt | Date | 建議存為 Date,便于范圍查詢。 |
| sentAt | Date | 首次發(fā)送時間。 |
| expireAt | Date | 可選。 |
| status | String | 見 2.4,必建索引。由「發(fā)送/重推」直接更新,或由 event_consumptions 聚合后異步/單獨更新,見 2.6。 |
| statusAt | Date | 最近一次狀態(tài)變更時間(發(fā)送、重推或匯總消費狀態(tài)時更新)。 |
| retryCount | Integer | 重試次數(shù)(含首次發(fā)送為 0,每次重推 +1)。 |
| lastSentAt | Date | 最近一次投遞到 MQ 的時間(含重推)。 |
| createdAt | Date | 文檔創(chuàng)建時間。 |
| updatedAt | Date | 文檔最后更新時間。 |
說明:消費反饋不寫入 events 文檔,統(tǒng)一寫入 event_consumptions 表并通過 eventId 關聯(lián),避免多消費者并發(fā)更新 events 導致沖突;events.status 的消費相關取值(CONSUMED/PARTIAL/FAILED)由 event_consumptions 聚合得到。
2.3 Topic–消費者配置表(topic_consumers)
用于維護 topic 與消費者的關聯(lián),支持動態(tài)配置:新增/下線某 topic 的消費者時只需改配置表(或通過 man 管理界面),無需發(fā)版。事件創(chuàng)建時會讀取該表并初始化 event_consumptions,見 3.2。
topic_consumers 文檔結構:
| 字段 | 類型 | 說明 |
|---|---|---|
| _id | ObjectId | MongoDB 主鍵。 |
| topic | String | 事件 topic,與 events.topic 一致,如 order.purchased、trade.payment.success。 |
| consumerId | String | 消費者唯一標識,如 member-service、message-service。 |
| enabled | Boolean | 是否啟用;僅 enabled=true 的配置會在事件創(chuàng)建時參與初始化。 |
| sortOrder | Integer | 可選,展示或初始化順序。 |
| description | String | 可選,消費者說明。 |
| createdAt | Date | 配置創(chuàng)建時間。 |
| updatedAt | Date | 配置最后更新時間。 |
唯一約束:(topic, consumerId) 唯一,同一 topic 下同一消費者只保留一條配置。
使用方式:
配置維護:通過 man 或直接寫庫增刪改 topic_consumers;例如為
order.purchased配置member-service、message-service,新事件創(chuàng)建時即會為這兩個消費者預插 event_consumptions 記錄。事件創(chuàng)建時初始化反饋:寫入 events 后,根據(jù)該事件的 topic 查詢 topic_consumers(enabled=true),為每個 consumerId 插入一條 event_consumptions 記錄:eventId=當前事件、consumerId=配置項、success=null(表示待消費)、consumedAt 為空、attemptNo=0(表示初始化)。后續(xù)消費者每次回調均新增一條 event_consumptions 記錄(attemptNo 遞增),不覆蓋舊記錄,便于保留多次重試/回調歷史、方便問題排查;匯總「當前狀態(tài)」時按 (eventId, consumerId) 取最新一條(attemptNo 最大或 consumedAt 最新)參與計算。管理端可區(qū)分「預期消費者數(shù)」與「已反饋數(shù)」,并展示「2/2 已消費」或「1/2 已消費」等,且可查看每個消費者的多次回調明細。
若某 topic 在 topic_consumers 中無配置(或全為 enabled=false):記錄 ERROR 日志、不拋異常、終止發(fā)送(不寫入 events、不發(fā)布 MQ),防止未配置消費者的 topic 被發(fā)出,且不影響調用方主流程,見 3.2。
2.4 狀態(tài)(status)枚舉
| 取值 | 含義 | 可遷轉方向 |
|---|---|---|
| PENDING | 已創(chuàng)建未發(fā)送(先落庫后發(fā) MQ 時,寫入 events 后、發(fā)布 MQ 前為該狀態(tài))。 | → SENT |
| SENT | 已發(fā)送到 RabbitMQ,尚未確認消費結果。 | → CONSUMED, PARTIAL, FAILED, RETRYING |
| CONSUMED | 已成功消費并確認(單消費者全部成功,或多消費者全部成功)。 | — |
| PARTIAL | 多消費者場景下,部分成功、部分失敗。 | → CONSUMED, FAILED, RETRYING |
| FAILED | 消費失敗或發(fā)送失敗,不再自動重試。 | → RETRYING(人工/管理端重推) |
| RETRYING | 正在重試(如已再次投遞到 MQ)。 | → SENT, CONSUMED, PARTIAL, FAILED |
| EXPIRED | 已過期(若有 expireAt)。 | — |
狀態(tài)變更時更新 status、statusAt;與消費相關的狀態(tài)(CONSUMED/PARTIAL/FAILED)不通過對 events 的 read-modify-write 直接改,而是由 event_consumptions 聚合后更新(見 2.5、2.6),從而避免并發(fā)寫沖突。
2.5 消費反饋表(event_consumptions)
消費反饋不入 events 文檔,全部寫入獨立集合 event_consumptions,通過 eventId 與 events 關聯(lián)。不設 (eventId, consumerId) 唯一約束:同一消費者對同一事件的多次回調各保留一條記錄,便于問題排查(如重試 3 次、3 條錯誤信息均可查看)。
每個消費者回調時僅插入新記錄,不更新已有記錄,不存在多線程寫同一文檔的并發(fā)問題。
單消費者與多消費者模型統(tǒng)一:同一 (eventId, consumerId) 可有多條記錄(初始化 1 條 + 每次回調 1 條),匯總「當前狀態(tài)」時按 (eventId, consumerId) 取最新一條參與計算。
event_consumptions 文檔結構:
| 字段 | 類型 | 說明 |
|---|---|---|
| _id | ObjectId | MongoDB 主鍵。 |
| eventId | String | 關聯(lián) events.eventId,必建索引。 |
| consumerId | String | 消費者唯一標識,如 member-service、message-service。單消費者場景可固定為同一值(如 default)。 |
| attemptNo | Integer | 同一 (eventId, consumerId) 下的序號:0 表示事件創(chuàng)建時按 topic_consumers 初始化的「待消費」記錄;1, 2, 3… 表示第 1、2、3 次回調。用于排序與取「最新一次」反饋。 |
| success | Boolean / null | 該次是否消費成功;null 僅用于 attemptNo=0 的初始化記錄(待消費)。 |
| consumedAt | Date | 該次消費完成時間;初始化記錄可為 null。 |
| errorMessage | String | 失敗時錯誤信息。 |
| errorCode | String | 可選,業(yè)務錯誤碼。 |
| createdAt | Date | 記錄創(chuàng)建時間。 |
無唯一約束。記錄來源:① 事件創(chuàng)建時按 topic_consumers 為每個 consumerId 插入一條(attemptNo=0、success=null);② 消費者每次回調均插入一條新記錄,attemptNo 為該 (eventId, consumerId) 下當前最大值+1。同一消費者多次重試/回調會形成多行,管理端可按時間或 attemptNo 查看完整歷史,便于排查。
與 events.status 的匯總關系: 按 eventId 查出所有 event_consumptions 后,先按 (eventId, consumerId) 取最新一條(attemptNo 最大,或 consumedAt 非空中取最大);再按這些「當前」記錄的 success 匯總:全部 success=true 且無待消費(null) → CONSUMED;存在 success=false → PARTIAL 或 FAILED;存在 success=null(僅初始化記錄)→ 仍可視為 SENT(等待中)或展示「已反饋數(shù)/預期數(shù)」。匯總結果不在消費回調路徑上直接寫回 events,而是由 2.6 的異步回寫(專用 topic 串行)回寫 events。
2.6 狀態(tài)更新與并發(fā)(events.status)
問題:若多個消費者同時回調,都在 events 上做「讀當前 status/consumptions → 計算新 status → 寫回」,會產(chǎn)生競態(tài)與覆蓋,且 events 文檔會因內嵌數(shù)組頻繁變更而膨脹。
做法:
消費反饋只寫 event_consumptions,且每次回調插入新記錄 各消費者(或 bus 代寫)每次回調向 event_consumptions 插入一條新記錄(attemptNo 遞增),不更新已有記錄,不修改 events。
events 上僅允許兩類更新,且互不疊加
發(fā)送/重推側:僅更新與「發(fā)送」相關的字段(如 status=PENDING→SENT、retryCount、lastSentAt、statusAt)。由 bus 單點或按 eventId 串行化后執(zhí)行,避免并發(fā)。
消費匯總側:僅更新「由消費結果推導出的 status」(CONSUMED/PARTIAL/FAILED)及 statusAt。不在此文檔上做數(shù)組 push 等復雜更新。
消費匯總采用「異步回寫」實現(xiàn),使用獨立 topic/隊列串行處理
實現(xiàn)方式:消費反饋寫入 event_consumptions 后,向消費匯總專用 topic(隊列)投遞一條輕量消息(如 payload 僅含 eventId),由該 topic 的單一消費者串行消費:根據(jù) eventId 查詢 event_consumptions、按 (eventId, consumerId) 取最新一條并匯總出 status(CONSUMED/PARTIAL/FAILED 或等待中),再對 events 做一次 status + statusAt 的更新。
專用 topic:與業(yè)務事件 topic 區(qū)分,由項目約定命名(如
bus.consumption-rollup、internal.status-rollup),單獨隊列、單獨消費者,不與業(yè)務消費混用。串行保證:該 topic 僅部署一個消費者(或同一消費者組內單實例),保證「匯總回寫」在同一時刻只處理一條消息,對 events 的更新不并發(fā),同一 eventId 的匯總只由一處執(zhí)行。
流程簡述:寫入 event_consumptions → 發(fā)送 eventId 到消費匯總專用 topic → 匯總消費者消費 → 查 event_consumptions 聚合 → 更新 events.status、statusAt。
2.7 索引建議
events:
db.events.createIndex({ eventId: 1 }, { unique: true })
db.events.createIndex({ status: 1, occurredAt: -1 })
db.events.createIndex({ topic: 1, occurredAt: -1 })
db.events.createIndex({ traceId: 1, occurredAt: 1 })
db.events.createIndex({ "initiator.service": 1, occurredAt: -1 })
db.events.createIndex({ parentEventId: 1 })
// db.events.createIndex({ expireAt: 1 }, { expireAfterSeconds: 0 })topic_consumers(按 topic 查配置,供事件創(chuàng)建時初始化 event_consumptions):
db.topic_consumers.createIndex({ topic: 1, consumerId: 1 }, { unique: true })
db.topic_consumers.createIndex({ topic: 1, enabled: 1 })event_consumptions(與 events 關聯(lián)、按消費者查詢、按 (eventId, consumerId) 取最新一條;無唯一約束,同一消費者多次回調多行):
db.event_consumptions.createIndex({ eventId: 1 })
db.event_consumptions.createIndex({ eventId: 1, consumerId: 1, attemptNo: -1 }) ?// 取某消費者在某事件下的最新反饋
db.event_consumptions.createIndex({ eventId: 1, consumerId: 1, consumedAt: -1 }) // 或按時間取最新
db.event_consumptions.createIndex({ consumerId: 1, consumedAt: -1 })3. 事件的發(fā)送
3.1 發(fā)送時機與職責
發(fā)送由 bus 完成:將事件發(fā)布到 RabbitMQ,并在 MongoDB 中寫入/更新事件記錄。
采用「先落庫后發(fā) MQ」:先以 PENDING 寫入 MongoDB(并初始化 event_consumptions),再發(fā)布到 RabbitMQ,成功后更新為 SENT。便于「至少存一次」、重推時從庫中取完整 envelope,且與消費反饋、重推流程一致。
3.2 發(fā)送流程(簡要)
生成/校驗 envelope:生成
eventId(若未提供)、occurredAt、sentAt,校驗topic、payload、initiator等。校驗 topic_consumers 配置:根據(jù)本事件的 topic 查詢 topic_consumers(enabled=true)。若無任何啟用中的消費者配置,則打印 ERROR 日志(如記錄 topic、eventId 等)、不拋異常,終止本次發(fā)送(不寫入 events、不發(fā)布 MQ),直接返回,不影響調用方主流程;調用方可在配置 topic_consumers 后重試或通過其他途徑處理。
寫入 MongoDB(events + 初始化消費反饋):
插入 events 文檔,status = PENDING。
初始化 event_consumptions:為步驟 2 查到的每個 consumerId 插入一條 event_consumptions 記錄:eventId=當前事件、consumerId=配置項、attemptNo=0、success=null、consumedAt=null(表示待消費)。
發(fā)布到 RabbitMQ:按
topic路由到對應 exchange/queue,消息體為序列化后的 envelope(或僅 payload + 必要元數(shù)據(jù),由實現(xiàn)約定)。更新發(fā)送狀態(tài):發(fā)布成功后,更新 events 的 status = SENT,寫入 sentAt、lastSentAt,retryCount 保持 0(首次)。
發(fā)送失?。ㄈ?MQ 不可用)時:已落庫的 events 保持 PENDING 或置為 FAILED,并記錄錯誤信息,便于 man 重推;已初始化的 event_consumptions 記錄保留,便于管理端展示「預期消費者」與后續(xù)真實反饋對比。
3.3 與「關聯(lián)」的關系
發(fā)送時若調用方傳入 traceId/spanId/parentEventId,應原樣寫入 envelope 和 MongoDB,用于后續(xù)「事件調用鏈路」展示。
initiator 在發(fā)送時由調用方傳入或由 bus 從當前上下文(服務名、請求 ID)補全,用于「事件發(fā)起方」展示。
4. 事件的保存
4.1 保存發(fā)生的時機
| 時機 | 寫入位置 | 內容 |
|---|---|---|
| 發(fā)送前/后 | events | 完整 envelope + status(PENDING/SENT)、sentAt、lastSentAt、retryCount、createdAt/updatedAt。 |
| 消費確認后 | event_consumptions | 插入一條新記錄(eventId、consumerId、attemptNo 遞增、success、consumedAt、errorMessage 等),不覆蓋已有記錄,便于保留多次回調歷史。不直接改 events 文檔。隨后向消費匯總專用 topic 投遞 eventId,由該 topic 的異步回寫消費者串行聚合后更新 events.status(見 2.6)。 |
| 重推時 | events | 僅更新 retryCount、lastSentAt、status(→ RETRYING/SENT)、statusAt。 |
| 人工/定時過期 | events | 更新 status = EXPIRED、statusAt。 |
4.2 冪等與唯一性
events:以 eventId 唯一標識一條事件;eventId 唯一索引,插入時若已存在則按策略忽略或更新發(fā)送相關字段,不在 events 上做消費反饋的 read-modify-write。
event_consumptions:無 (eventId, consumerId) 唯一約束;同一消費者對同一事件每次回調插入一行(attemptNo 遞增),歷史可查、便于排查,且無并發(fā)更新 events 的問題。
發(fā)送到 RabbitMQ 時,可根據(jù) eventId 做去重;消費端按 eventId(及 consumerId)做業(yè)務冪等。
4.3 誰負責寫 MongoDB
bus:發(fā)送時寫 events,并據(jù) topic_consumers 初始化 event_consumptions(見 3.2);消費端回調或 MQ 確認時只寫 event_consumptions,并投遞 eventId 到消費匯總專用 topic。events 的消費匯總 status 按 2.6 由該 topic 的異步回寫消費者串行聚合后回寫。
man:讀 events 與 event_consumptions(列表/詳情可做關聯(lián)或聚合);重推時只更新 events 的 retryCount、lastSentAt、status 等,由 bus 或 man 單點執(zhí)行。
5. 重試與重推
5.1 概念區(qū)分
自動重試:消費失敗后由 MQ 或 bus 自動重新投遞(如 RabbitMQ 重入隊、DLQ 再投遞),次數(shù)與策略可配置。
人工/管理端重推:在 man 上對某條事件(如 FAILED 或 SENT 超時未消費)執(zhí)行「重推」,再次將該事件投遞到 RabbitMQ。
本文檔側重 管理端重推;自動重試策略可與 RabbitMQ 的 prefetch、nack、DLQ 等配合,由 bus 配置實現(xiàn)。
5.2 重推條件(建議)
status 為 FAILED 或 RETRYING;或 SENT 且超過一定時間未變?yōu)?CONSUMED(超時未消費)。
若存在 expireAt,則重推時檢查未過期再推。
可選:限制單條事件最大重推次數(shù),超過則不再允許重推。
5.3 重推流程
man 接收重推請求(如按 eventId)。
man 從 MongoDB 根據(jù) eventId 查出完整事件文檔(envelope)。
man 調用 bus 的重推接口,或直接向 RabbitMQ 發(fā)布該事件(若 man 具備 MQ 客戶端權限);推薦通過 bus 統(tǒng)一發(fā)布,便于審計與計數(shù)。
bus(或 man)將事件再次發(fā)布到 RabbitMQ,并更新 MongoDB:
retryCount += 1;
lastSentAt = 當前時間;
status = RETRYING 或 SENT。
消費端再次消費后,按 4.1 向 event_consumptions 插入該消費者的一條新記錄(attemptNo 遞增),并投遞 eventId 到消費匯總專用 topic;events 的匯總 status 由該 topic 的異步回寫消費者串行回寫(見 2.6)。
5.4 冪等與重復消費
重推使用同一 eventId,消費端應基于 eventId 做冪等處理(如先查已處理再執(zhí)行業(yè)務),避免重復扣款、重復發(fā)通知等。
存儲側同一 eventId 只對應一條文檔,重推只更新該文檔,不新增文檔。
6. 關聯(lián):發(fā)起方與事件調用鏈路
6.1 事件發(fā)起方
initiator 即發(fā)起方信息,在發(fā)送時寫入 envelope 與 MongoDB。
man 列表/詳情可展示:initiator.service、initiator.operation、initiator.userId、initiator.clientRequestId,并支持按 initiator.service 等篩選。
若發(fā)送時未傳 initiator,可由 bus 用默認值(如當前服務名、當前請求 ID)補全,保證「總能展示發(fā)起方」。
6.2 事件調用鏈路
traceId:同一次業(yè)務觸發(fā)的多個事件/步驟共享同一 traceId,用于「按鏈路聚合」。
parentEventId:若事件 A 的消費邏輯里又發(fā)送了事件 B,則 B 的 parentEventId = A.eventId,形成父子關系。
spanId:同一 trace 下不同節(jié)點的唯一標識,便于在界面上畫成樹或時間線。
展示方式建議:
按 traceId 查詢:列出該 traceId 下所有事件,按 occurredAt 或 span 順序排序,展示為時間線或樹(父 → 子)。
按 eventId 查詳情:展示該事件 + 其 parent(parentEventId)及子事件(parentEventId = 當前 eventId),即「上下游」各一層或多層。
存儲上不強制 event_links 表;用 events 的 traceId、parentEventId、eventId 即可通過查詢實現(xiàn)「鏈路」與「樹形」展示。
6.3 關聯(lián)關系小結
| 關系 | 存儲字段 | 用途 |
|---|---|---|
| 發(fā)起方 | initiator | 誰/哪個服務/哪個操作發(fā)起了事件 |
| 同鏈路 | traceId | 哪些事件屬于同一次業(yè)務過程 |
| 父子事件 | parentEventId ↔ eventId | 誰觸發(fā)了誰(下游由上游消費觸發(fā)) |
| 同鏈路順序 | spanId + occurredAt | 鏈路內節(jié)點順序與并行關系 |
7. 典型場景:一源多消費者(購買成功 → 會員積分 + 消息通知)
7.1 場景描述
業(yè)務:用戶 A 購買商品成功。
期望:同一筆購買需要觸發(fā)多個下游動作且互不阻塞:
會員服務:根據(jù)訂單金額增加積分;
消息服務:發(fā)送通知(站內信/短信/推送等)。
要求事件設計能兼容「一個業(yè)務事件、多個獨立消費者」:發(fā)一次事件,會員與消息各自消費、各自可追溯、各自可重試,且能在管理端看到完整鏈路(誰發(fā)起的、誰消費了、誰又觸發(fā)了子事件)。
7.2 事件設計(兼容多消費者)
原則:一條業(yè)務事實只對應一條系統(tǒng)事件(一個 eventId),payload 攜帶所有消費者都可能需要的公共信息;各消費者按需取用,不要求 payload 為某單一消費者定制。
topic:用業(yè)務含義命名,與「有多少個消費者」無關。例如:
order.purchased或trade.payment.success。payload:包含訂單、用戶、金額等公共字段,供會員服務算積分、消息服務填模板。建議至少包含:
orderId、userId、amount、occurredAt(或訂單完成時間);可選:
channel、productIds、quantity等,按業(yè)務約定。
initiator:填寫發(fā)起方,如
service: "order-service",operation: "confirmOrder",userId: "user-A",便于在 man 中看到「用戶 A 購買成功」是誰發(fā)起的。
示例 payload(僅示例,可按業(yè)務擴展):
{
?"orderId": "ORD-2024-002",
?"userId": "user-A",
?"amount": 199.00,
?"currency": "CNY",
?"occurredAt": "2024-02-28T10:05:00Z",
?"channel": "app"
}會員服務只關心 userId、amount 等;消息服務只關心 userId、orderId、amount 等;兩者共用同一 payload,無需為每個消費者發(fā)不同事件。
7.3 消息路由(RabbitMQ 一事件多隊列)
同一事件需要被會員服務和消息服務各自消費一次,且互不影響(一個失敗不影響另一個)。推薦兩種方式二選一:
| 方式 | 說明 |
|---|---|
| Fanout 交換機 | 訂單服務將事件發(fā)到 fanout exchange,該 exchange 綁定隊列 member-service.order.purchased、message-service.order.purchased;每個隊列一份消息,各消費者獨立消費、獨立 ack。 |
| Topic 交換機 + 多 binding | 使用同一 topic(如 order.purchased),兩個隊列用相同或不同 binding key 訂閱該 topic;broker 將消息復制到每個隊列。 |
無論哪種方式,存儲側只存一條事件文檔(一個 eventId):該事件被投遞到多個隊列,但仍是「同一條事件」;消費反饋通過 event_consumptions 表按消費者區(qū)分;預期消費者由 topic_consumers 配置并在事件創(chuàng)建時初始化(見 2.3、2.5、7.4)。
7.4 存儲與狀態(tài)(多消費者)
events 表仍是一條文檔對應一條業(yè)務事件(一次「用戶 A 購買成功」對應一個 eventId);不在 events 文檔內嵌 consumptions,避免多消費者并發(fā)更新同一文檔。
event_consumptions 表按 (eventId, consumerId) 可有多條記錄(無唯一約束),例如該 eventId 下:
| eventId | consumerId | attemptNo | success | consumedAt |
|---|---|---|---|---|
| evt-001 | member-service | 0 | null | null(初始化) |
| evt-001 | member-service | 1 | true | 2024-02-28T10:05:02Z |
| evt-001 | message-service | 0 | null | null(初始化) |
| evt-001 | message-service | 1 | false | 2024-02-28T10:05:03Z |
| evt-001 | message-service | 2 | true | 2024-02-28T10:05:10Z(重試成功) |
匯總時按 (eventId, consumerId) 取 attemptNo 最大(或 consumedAt 最新)的一條參與計算;管理端可查看每個消費者的多次回調歷史,便于排查重試與錯誤信息。
events.status 的匯總規(guī)則(見 2.4、2.6):按 (eventId, consumerId) 取最新一條后聚合 —— 全部 success=true 且無待消費(null) → CONSUMED;存在 success=false → PARTIAL 或 FAILED;存在 success=null 可視為等待中。匯總結果由消費匯總專用 topic 的異步回寫消費者串行處理并回寫 events,不在消費回調路徑上直接改 events。
這樣在管理端可以:看「這條購買事件」的總狀態(tài)(來自 events 或實時聚合),并關聯(lián) event_consumptions 展開「會員服務已成功、消息服務失敗」等明細及每次回調的 attemptNo/時間/錯誤信息,便于按消費者排查與重試。
7.5 子事件與鏈路(同一 traceId、parentEventId)
會員服務、消息服務在處理完該事件后,若會再產(chǎn)生下游事件(如「積分已增加」「通知已發(fā)送」),應作為子事件發(fā)送,以便在 man 中形成完整鏈路:
父事件:
order.purchased(eventId = evt-purchase-001,traceId = trace-001)。子事件 1:會員服務發(fā)送
points.added,攜帶parentEventId = evt-purchase-001,traceId = trace-001(與父一致),payload 可含 userId、points、sourceEventId 等。子事件 2:消息服務發(fā)送
notification.sent,攜帶parentEventId = evt-purchase-001,traceId = trace-001,payload 可含 userId、channel、templateId 等。
這樣在「事件調用鏈路」中可按 traceId 展示為: order.purchased → points.added、order.purchased → notification.sent(兩條并行分支),且能看出都源自同一次「用戶 A 購買成功」。
7.6 重推策略(按需按消費者)
整事件重推:對 eventId 重推時,事件會再次投遞到所有訂閱了該 topic 的隊列,會員和消息都會再收一次;各消費者需按 eventId(+ 自身 consumerId)做冪等,避免重復加積分、重復發(fā)通知。
按消費者重推(可選):若 man 支持「僅對某消費者重推」,則需總線或 MQ 支持按隊列/消費者定向投遞(例如單獨發(fā)到 message-service 的隊列);存儲上該消費者再次回調時會新插入一條 event_consumptions 記錄(attemptNo 遞增),歷史可查。當前文檔以「整事件重推」為主,按消費者重推可作為擴展實現(xiàn)。
7.7 小結(兼容性要點)
| 維度 | 設計要點 |
|---|---|
| 事件條數(shù) | 一次業(yè)務事實 = 一條事件(一個 eventId),不按消費者拆成多條。 |
| Payload | 設計成多消費者可共用的公共信息,便于擴展更多下游(如風控、統(tǒng)計)。 |
| 路由 | RabbitMQ 一事件多隊列(fanout 或 topic 多 binding),每個消費者獨立 ack。 |
| 存儲 | events 僅存事件與發(fā)送態(tài);消費反饋存 event_consumptions 表,通過 eventId 關聯(lián);匯總 status 采用異步回寫,經(jīng)專用 topic 串行更新 events(見 2.6),避免并發(fā)更新。 |
| 鏈路 | 子事件用同一 traceId + parentEventId,便于「購買 → 積分、通知」同鏈展示。 |
| 重推與冪等 | 重推同一 eventId,各消費者按 eventId(+ consumerId)做冪等。 |
8. 流程總覽(Mermaid)
以下為「發(fā)送 → 存儲 → 消費 → 狀態(tài)更新」與「重推」的簡化流程。
man消費者RabbitMQMongoDBbus調用方man消費者RabbitMQMongoDBbus調用方發(fā)送與保存重推發(fā)布事件(envelope)寫入 events (PENDING)發(fā)布消息ack更新 SENT, sentAt, lastSentAt投遞業(yè)務處理消費結果(成功/失敗)插入 event_consumptions 新記錄(eventId, consumerId, attemptNo),不直接改 events按 eventId 查詢重推(eventId)retryCount+1, lastSentAt, status再次發(fā)布再次投遞
文檔版本:v0.4 狀態(tài):設計說明,供實現(xiàn) bus / man 時參考。若后續(xù)有字段或流程變更,請同步更新本文檔與 需求文檔。
一源多消費者場景(如:購買成功 → 會員積分 + 消息通知)見 第 7 節(jié)。
消費反饋存 event_consumptions 獨立表,與 events 通過 eventId 關聯(lián),避免狀態(tài)更新并發(fā)問題,見 2.1、2.5、2.6。
event_consumptions 無 (eventId, consumerId) 唯一約束:同一消費者多次回調各插入一條記錄(attemptNo 遞增),歷史可查、便于問題排查;匯總時按 (eventId, consumerId) 取最新一條,見 2.5。
Topic–消費者配置表 topic_consumers:維護 topic 與消費者關聯(lián),支持動態(tài)配置;事件創(chuàng)建時據(jù)該表初始化 event_consumptions(預期消費者),見 2.3、3.2。
到此這篇關于Java微服務事件與存儲設計說明的文章就介紹到這了,更多相關Java微服務與存儲設計內容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關文章希望大家以后多多支持腳本之家!
相關文章
淺析java中String類型中“==”與“equal”的區(qū)別
這篇文章主要介紹了淺析java中String類型中“==”與“equal”的區(qū)別,本文給大家介紹的非常詳細,對大家的學習或工作具有一定的參考借鑒價值,需要的朋友可以參考下2020-08-08
SpringBoot監(jiān)聽應用程序啟動的生命周期事件的四種方法
在 Spring Boot 中,監(jiān)聽應用程序啟動的生命周期事件有多種方法,本文給大家就介紹了四種監(jiān)聽應用程序啟動的生命周期事件的方法,并通過代碼示例講解的非常詳細,具有一定的參考價值,需要的朋友可以參考下2024-07-07
解析Mybatis的insert方法返回數(shù)字-2147482646的解決
這篇文章主要介紹了解析Mybatis的insert方法返回數(shù)字-2147482646的解決,文中通過示例代碼介紹的非常詳細,對大家的學習或者工作具有一定的參考學習價值,需要的朋友們下面隨著小編來一起學習學習吧2021-04-04
WIN7系統(tǒng)JavaEE(java)環(huán)境配置教程(一)
這篇文章主要介紹了WIN7系統(tǒng)JavaEE(java+tomcat7+Eclipse)環(huán)境配置教程,本文重點在于java配置,感興趣的小伙伴們可以參考一下2016-06-06

