最新国产好看的视频,伊人天堂AV在线,国产Aaaaaa视频,蜜臀视频在线观看一区,人妻av色图,密臀久久久精品影片,青青视频免费观看毛片,久草在线观看视,国产三级精品色情在线

Java微服務事件與存儲設計說明(含代碼示例)

 更新時間:2026年06月30日 10:08:32   作者:小莫分享  
ava微服務架構的設計與優(yōu)化隨著企業(yè)系統(tǒng)的規(guī)模不斷擴大,傳統(tǒng)的單體應用架構已經(jīng)無法滿足高效開發(fā)、擴展性、容錯性等需求,這篇文章主要介紹了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)。

字段類型必填說明
eventIdString全局唯一事件 ID(如 UUID),用于去重、重試、關聯(lián)。
traceIdString推薦分布式鏈路 ID,同一次業(yè)務觸發(fā)的多條事件/調用共享,用于串聯(lián)「事件調用鏈路」。
spanIdString可選當前環(huán)節(jié)的 span ID,與 traceId 一起構成鏈路節(jié)點。
parentEventIdString可選若本事件由另一事件的消費所觸發(fā),則填上游事件的 eventId,用于父子關聯(lián)。
topicString邏輯主題,對應 RabbitMQ 的 routing key 或交換機路由,如 order.created、payment.notify
payloadObject / String / Binary業(yè)務載荷,JSON 對象或序列化字符串;消費端按 topic 約定反序列化。
payloadTypeString可選載荷類型提示,如 application/json、application/cloudevents+json
initiatorObject推薦發(fā)起方信息,見下。
occurredAtString (ISO-8601)事件發(fā)生時間(業(yè)務發(fā)生時間或首次發(fā)送時間)。
sentAtString (ISO-8601)可選首次發(fā)送到 MQ 的時間,由 bus 在發(fā)送時寫入。
expireAtString (ISO-8601)可選業(yè)務過期時間,超過后可不再重試或不再展示。

initiator(發(fā)起方) 建議結構:

字段類型說明
serviceString發(fā)起服務名,如 order-servicebus。
operationString操作/接口,如 createOrderPOST /api/orders。
userIdString若可識別,填用戶 ID。
clientRequestIdString客戶端請求 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_consumersTopic–消費者配置表:維護每個 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)字段:

字段類型說明
_idObjectIdMongoDB 主鍵;也可用 eventId 作為 _id 便于冪等寫入。
eventIdString同 1.2,唯一,建議建唯一索引。
traceIdString同 1.2,建索引便于按鏈路查詢。
spanIdString同 1.2。
parentEventIdString同 1.2,建索引便于查「某事件的子事件」。
topicString同 1.2,必建索引(列表/篩選)。
payloadObject / String同 1.2。
payloadTypeString同 1.2。
initiatorObject同 1.2,可按 initiator.service 等建復合索引。
occurredAtDate建議存為 Date,便于范圍查詢。
sentAtDate首次發(fā)送時間。
expireAtDate可選。
statusString見 2.4,必建索引。由「發(fā)送/重推」直接更新,或由 event_consumptions 聚合后異步/單獨更新,見 2.6。
statusAtDate最近一次狀態(tài)變更時間(發(fā)送、重推或匯總消費狀態(tài)時更新)。
retryCountInteger重試次數(shù)(含首次發(fā)送為 0,每次重推 +1)。
lastSentAtDate最近一次投遞到 MQ 的時間(含重推)。
createdAtDate文檔創(chuàng)建時間。
updatedAtDate文檔最后更新時間。

說明:消費反饋不寫入 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 文檔結構

字段類型說明
_idObjectIdMongoDB 主鍵。
topicString事件 topic,與 events.topic 一致,如 order.purchased、trade.payment.success。
consumerIdString消費者唯一標識,如 member-service、message-service
enabledBoolean是否啟用;僅 enabled=true 的配置會在事件創(chuàng)建時參與初始化。
sortOrderInteger可選,展示或初始化順序。
descriptionString可選,消費者說明。
createdAtDate配置創(chuàng)建時間。
updatedAtDate配置最后更新時間。

唯一約束(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 文檔結構

字段類型說明
_idObjectIdMongoDB 主鍵。
eventIdString關聯(lián) events.eventId,必建索引。
consumerIdString消費者唯一標識,如 member-servicemessage-service。單消費者場景可固定為同一值(如 default)。
attemptNoInteger同一 (eventId, consumerId) 下的序號:0 表示事件創(chuàng)建時按 topic_consumers 初始化的「待消費」記錄;1, 2, 3… 表示第 1、2、3 次回調。用于排序與取「最新一次」反饋。
successBoolean / null該次是否消費成功;null 僅用于 attemptNo=0 的初始化記錄(待消費)。
consumedAtDate該次消費完成時間;初始化記錄可為 null。
errorMessageString失敗時錯誤信息。
errorCodeString可選,業(yè)務錯誤碼。
createdAtDate記錄創(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ù)組頻繁變更而膨脹。

做法

  1. 消費反饋只寫 event_consumptions,且每次回調插入新記錄 各消費者(或 bus 代寫)每次回調向 event_consumptions 插入一條新記錄(attemptNo 遞增),不更新已有記錄,不修改 events。

  2. events 上僅允許兩類更新,且互不疊加

    • 發(fā)送/重推側:僅更新與「發(fā)送」相關的字段(如 status=PENDING→SENT、retryCount、lastSentAt、statusAt)。由 bus 單點或按 eventId 串行化后執(zhí)行,避免并發(fā)。

    • 消費匯總側:僅更新「由消費結果推導出的 status」(CONSUMED/PARTIAL/FAILED)及 statusAt。不在此文檔上做數(shù)組 push 等復雜更新。

  3. 消費匯總采用「異步回寫」實現(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ā)送流程(簡要)

  1. 生成/校驗 envelope:生成 eventId(若未提供)、occurredAt、sentAt,校驗 topic、payloadinitiator 等。

  2. 校驗 topic_consumers 配置:根據(jù)本事件的 topic 查詢 topic_consumers(enabled=true)。若無任何啟用中的消費者配置,則打印 ERROR 日志(如記錄 topic、eventId 等)、不拋異常,終止本次發(fā)送(不寫入 events、不發(fā)布 MQ),直接返回,不影響調用方主流程;調用方可在配置 topic_consumers 后重試或通過其他途徑處理。

  3. 寫入 MongoDB(events + 初始化消費反饋)

    • 插入 events 文檔,status = PENDING

    • 初始化 event_consumptions:為步驟 2 查到的每個 consumerId 插入一條 event_consumptions 記錄:eventId=當前事件、consumerId=配置項、attemptNo=0success=null、consumedAt=null(表示待消費)。

  4. 發(fā)布到 RabbitMQ:按 topic 路由到對應 exchange/queue,消息體為序列化后的 envelope(或僅 payload + 必要元數(shù)據(jù),由實現(xiàn)約定)。

  5. 更新發(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:讀 eventsevent_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 為 FAILEDRETRYING;或 SENT 且超過一定時間未變?yōu)?CONSUMED(超時未消費)。

  • 若存在 expireAt,則重推時檢查未過期再推。

  • 可選:限制單條事件最大重推次數(shù),超過則不再允許重推。

5.3 重推流程

  1. man 接收重推請求(如按 eventId)。

  2. man 從 MongoDB 根據(jù) eventId 查出完整事件文檔(envelope)。

  3. man 調用 bus 的重推接口,或直接向 RabbitMQ 發(fā)布該事件(若 man 具備 MQ 客戶端權限);推薦通過 bus 統(tǒng)一發(fā)布,便于審計與計數(shù)。

  4. bus(或 man)將事件再次發(fā)布到 RabbitMQ,并更新 MongoDB:

    • retryCount += 1;

    • lastSentAt = 當前時間;

    • status = RETRYING 或 SENT。

  5. 消費端再次消費后,按 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.purchasedtrade.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 下:

eventIdconsumerIdattemptNosuccessconsumedAt
evt-001member-service0nullnull(初始化)
evt-001member-service1true2024-02-28T10:05:02Z
evt-001message-service0nullnull(初始化)
evt-001message-service1false2024-02-28T10:05:03Z
evt-001message-service2true2024-02-28T10:05:10Z(重試成功)

匯總時按 (eventId, consumerId) 取 attemptNo 最大(或 consumedAt 最新)的一條參與計算;管理端可查看每個消費者的多次回調歷史,便于排查重試與錯誤信息。

  • events.status 的匯總規(guī)則(見 2.4、2.6):按 (eventId, consumerId) 取最新一條后聚合 —— 全部 success=true 且無待消費(null) → CONSUMED;存在 success=false → PARTIALFAILED;存在 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.purchasedpoints.added、order.purchasednotification.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ū)別

    這篇文章主要介紹了淺析java中String類型中“==”與“equal”的區(qū)別,本文給大家介紹的非常詳細,對大家的學習或工作具有一定的參考借鑒價值,需要的朋友可以參考下
    2020-08-08
  • Java方法引用原理實例解析

    Java方法引用原理實例解析

    這篇文章主要介紹了Java方法引用的原理實例解析,文中通過示例代碼介紹的非常詳細,對大家的學習或者工作具有一定的參考學習價值,需要的朋友可以參考下
    2021-08-08
  • 詳解JAVA之運算符

    詳解JAVA之運算符

    這篇文章主要介紹了詳解Java中運算符以及相關的用法講解,一起跟著小編學習下吧,希望能夠給你帶來幫助
    2021-11-11
  • Java實現(xiàn)UTF-8編碼與解碼方式

    Java實現(xiàn)UTF-8編碼與解碼方式

    這篇文章主要介紹了Java實現(xiàn)UTF-8編碼與解碼方式,具有很好的參考價值,希望對大家有所幫助。如有錯誤或未考慮完全的地方,望不吝賜教
    2023-04-04
  • SpringBoot監(jiān)聽應用程序啟動的生命周期事件的四種方法

    SpringBoot監(jiān)聽應用程序啟動的生命周期事件的四種方法

    在 Spring Boot 中,監(jiān)聽應用程序啟動的生命周期事件有多種方法,本文給大家就介紹了四種監(jiān)聽應用程序啟動的生命周期事件的方法,并通過代碼示例講解的非常詳細,具有一定的參考價值,需要的朋友可以參考下
    2024-07-07
  • 解析Mybatis的insert方法返回數(shù)字-2147482646的解決

    解析Mybatis的insert方法返回數(shù)字-2147482646的解決

    這篇文章主要介紹了解析Mybatis的insert方法返回數(shù)字-2147482646的解決,文中通過示例代碼介紹的非常詳細,對大家的學習或者工作具有一定的參考學習價值,需要的朋友們下面隨著小編來一起學習學習吧
    2021-04-04
  • WIN7系統(tǒng)JavaEE(java)環(huán)境配置教程(一)

    WIN7系統(tǒng)JavaEE(java)環(huán)境配置教程(一)

    這篇文章主要介紹了WIN7系統(tǒng)JavaEE(java+tomcat7+Eclipse)環(huán)境配置教程,本文重點在于java配置,感興趣的小伙伴們可以參考一下
    2016-06-06
  • JAVA按字節(jié)讀取文件的簡單實例

    JAVA按字節(jié)讀取文件的簡單實例

    下面小編就為大家?guī)硪黄狫AVA按字節(jié)讀取文件的簡單實例。小編覺得挺不錯的,現(xiàn)在就分享給大家,也給大家做個參考。一起跟隨小編過來看看吧
    2016-08-08
  • Java實現(xiàn)控制小數(shù)精度的方法

    Java實現(xiàn)控制小數(shù)精度的方法

    這篇文章主要介紹了Java實現(xiàn)控制小數(shù)精度的方法,文中通過示例代碼介紹的非常詳細,對大家的學習或者工作具有一定的參考學習價值,需要的朋友們下面隨著小編來一起學習學習吧
    2021-01-01
  • 教你怎么用java一鍵自動生成數(shù)據(jù)庫文檔

    教你怎么用java一鍵自動生成數(shù)據(jù)庫文檔

    最近小編也在找這樣的插件,就是不想寫文檔了,浪費時間和心情啊,果然我找到一款比較好用,操作簡單不復雜.screw 是一個簡潔好用的數(shù)據(jù)庫表結構文檔的生成工具,支持 MySQL、Oracle、PostgreSQL 等主流的關系數(shù)據(jù)庫.需要的朋友可以參考下
    2021-05-05

最新評論

南涧| 惠安县| 安新县| 柯坪县| 丹棱县| 巫溪县| 咸阳市| 邹城市| 林周县| 东乌珠穆沁旗| 亳州市| 正宁县| 鹰潭市| 靖宇县| 葫芦岛市| 青冈县| 岳阳市| 河南省| 长葛市| 十堰市| 龙口市| 巨鹿县| 泰和县| 呼图壁县| 锡林郭勒盟| 凌源市| 屏山县| 滦平县| 黄浦区| 浏阳市| 故城县| 西宁市| 久治县| 东莞市| 马鞍山市| 阿勒泰市| 安义县| 章丘市| 裕民县| 田东县| 龙州县|