l

2022年7月12日 星期二

事件溯源(15):Apache Kafka可以當做Event Store嗎?

July 5 19:02~20:09

圖片擷取自維基百科

 

前言

Teddy剛開始接觸Event Sourcing時,就是使用EventStoreDB,後來在ezKanban中也支援用PostgreSQL當作Event Store。但是在YouTube聽演講或是看網路文章,有時會發現有人把Apache Kafka做為Event Store。前陣子Teddy在看Apache Pulsar的書,Pulsar是和Kafka類似的Event Broker軟體但感覺好像架構又設計的更好一些。當時Teddy就想:「可以拿Pulsar當Event Store嗎?」

Kafka與Pulsar在message-oriented system中被大量使用,它們具有可靠、高效、可拓展與可永久儲存訊息的能力,很自然地會讓人想把它們當成Event Store使用。Teddy在網路上沒有找到Pulsar是否可做為Event Store的相關討論,倒是看到一篇文章<Apache Kafka Is Not for Event Sourcing>。因為不少作Event Sourcing的人可能也都會有這樣的疑問,Teddy今天就介紹這篇文章中提到不適合的兩點理由。

 

***

載入狀態

在Kafka中,訊息放在Topic裡面,然後Topic通常以entity type做為分類。以ezKanban的Core Domain為例,如果用Kafka當Event Store,會有Boards, Workflows, Cards, Tags這四個Topic。假設要讀取出某張Card,需要從Cards Topic讀出所有卡片的所有訊息,然後依據card id過濾出所需卡片的相關訊息。雖然這樣做是可行的,但有點不切實際,可能會有效率低落的疑慮。

如果學習EventStoreDB的做法,一個Aggregate instance享有一個Topic,這麼一來就會有為數非常可觀的Topic在Kafka伺服器上面。姑且不論Kafka能不能支撐大量的Topic,對於subscriber而言,要如何接收的想要的訊息(領域事件)?舉個例子,想要知道所有的CardMoved事件,在EventStoreDB中,只需要從$et-CardMoved這個系統自動投影出來的event stream就可以拿到。但在Kafka的情況下,subscriber需要去註冊所有Card instance的Topic才能拿到所有的CardMoved事件,這有點窒礙難行。

***


一致性寫入

文章中提到第二點不適合的理由是Teddy在<事件溯源(7):樂觀鎖>介紹過樂觀鎖的問題。在併行處理的系統中,為了避免寫入衝突,通常會採用樂觀鎖的機制,在一般的關聯式資料庫或是EventStoreDB都支援。但依據該文章的說法,Kakfa的Topic寫入沒有支援樂觀鎖,因此也不合適做為Event Store使用。

 

***

使用時機

Kafka或Pulsar這種Event Broker不是說在Event Sourcing系統中都不能使用,只是不要拿來當作Event Store。如下圖所示,領域事件還是存在EventStoreDB或是其他類似的資料庫,至於要往外(下游,downstream)傳遞的領域事件,可以轉發至Kafka或Pulsar。

換句話說,把Kafka或Pulsar拿來當作跨Bounded Context或是跨服務的訊息傳遞,不要拿來做Event Store。

 

 

***

 

下集預告

這一集內容雖然比較簡短,但卻是很重要個觀念,因為選對Event Store當你上天堂,選錯Event Store會讓你住套房。下一集回頭介紹一組相關的基礎觀念:在最終一致性(eventual consistency)的情況下,事件必須滿足event ordering與at least once這兩個條件,以及event handler要如何達到idempotent。

***

友藏內心獨白:很想用一個工具打通關。

2022年7月11日 星期一

事件溯源(14):行為版本控制

July 05 21:59~23:12

▲圖1:行為無法控制

 

前言

上一集討論事件版本控制的議題,這一集要討論行為版本控制(Behavior Versioning)。這個問題Teddy也是讀了Gregory Young的《Versioning in an Event Sourced System》才發現,如果沒注意到程式的「行為版本」,在Event Sourcing系統中,可能會因為程式改版(領域事件不變)導致系統狀態錯誤。

 

***

問題範例

Gregory Young的在《Versioning in an Event Sourced System》書中有一個很簡單的例子來說明行為版本控制,在ezKanban中也有類似的例子但比較複雜一些,在這裡Teddy直接借用Gregory Young的例子。如圖2所示,銷售系統的POS Aggregate身上有sell()方法,它apply一個ItemSold領域事件,其中事件的最後一個參數是這筆銷售的金額小計。

 


▲圖2:sell method

 

圖3是ItemSold的event handler,它從領域事件拿出subTotal並將其乘上0.08以計算營業稅(tax是POS Aggregate身上的屬性);程式這樣寫看起來沒什麼問題。

 

▲圖3:計算稅金,1.0版程式

 

有一天政府把營業稅調高,從8%調整成10%,於是你把程式改成圖4。現在問題來了,當你重新載入POS之後,它會replay所有事件,然後假設原本有一個項目subTotal是100,tax採用1.0版程式計算出來的tax是8。但是程式改版之後,用2.0版程式計算出的tax變成了10。但是這筆交易產生的時間點,營業稅還是8%,所以不應該因為程式碼改變,就造成「竄改歷史」的情況。

 

▲圖4:計算稅金,2.0版程式

 

***


解決方法

如圖5所示,解決方法其實很簡單,就是把tax算好並當成領域事件內容的一部分,這樣就可以了。回到領域事件原始定義:「代表系統狀態改變」,所以領域事件內容應該「至少」要儲存「能夠代表本次狀態改變的所有資料。」在這個例子中,「稅率」是會隨著時間改變的數值,會影響tax的金額。因此ItemSold領域事件應該包含計算後的tax,或是只包含taxRate(稅率),之後再依據taxRate去計算tax也是可以。

 

▲圖5:修改後的版本,不會受到程式行為調整而改變聚合狀態

 

***

以上這個例子算是簡單的,在某些情況底下,如果Aggregate呼叫外部服務,也很可能會造成行為版本控制的問題,因而導致無法replay領域事件重現系統狀態。例如,向第三方金流API請款,如果replay領域事件會不會導致重複請款?這些都是要注意的細節。呼叫外部服務導致程式行為在replay變得不可決定(nondeterministic)的解決方法和上面計算稅金的例子類似,在產生領域事件時呼叫外部服務,並且把外部服務的回傳值儲存在領域事件上。如此可以達到確定性重播(deterministic replay)。總而言之,重點就是重建系統狀態所需的所有資料都要儲存在event stream裡面,如此每次replay才會出現相同的結果,如圖6所示。

 

▲圖6:將外部服務回傳結果存入領域事件中

 

Gregory Young在《Versioning in an Event Sourced System》書中還提到更多關於行為版本控制的細節,有興趣的鄉民請自行參考。

***

 

下集預告

在這一系列的文章中,Teddy使用EventStoreDB與PostgreSQL當作Event Store。早期採用Event Sourcing的開發人員有不少採用Apache Kafka當作Event Store。下一集要談是否合適把Apache Kafka當做Event Store?

***

友藏內心獨白:重建犯罪現場真的沒有那麼簡單啊。

2022年7月10日 星期日

事件溯源(13):事件版本控制

July 05 18:08~19:23;20:07~21:06

▲圖1:喵星人需要做版本控管的嗎?

 

前言

Event Sourcing因為把事件當作單一真實資料來源(single source of truth),因此事件本身就相當於是一種API。事件改版就好像API改版一樣,對於Event Sourcing系統來講是一個很重要的議題,如果不同版本的事件不相容,之後在replay事件的時候就可能會出問題,導致系統錯誤甚至是無法使用。這一集介紹事件版本控制(Event Versioning)的基本觀念與做法。

 

***

Schema on Write and Schema on Read

傳統關聯式資料庫當資料寫入資料表格時,必須知道這個table schema才可以成功寫入資料,這種方式稱為schema-on-write。需要事先定義資料格式,然後依據這個格式寫入資料庫。好處是在寫入時可以驗證資料格式的正確性,缺點則是因為資料格式在寫入時就已經固定,所以比較沒有彈性。

NoSQL資料庫(包含EventStoreDB)以及Apache Pulsar event broker則是採用另一種稱為schema-on-read方式,寫入時不檢查資料格式(資料可能以byte array的方式存在資料中),因此比較有彈性;至於資料的正確性則是交由讀取者來負責。

在Event Sourcing系統中,領域事件的格式都不一樣,因此幾乎不可能採用schema-on-write的方式來儲存領域事件,所以自然地採用schema-on-read。但問題來了,如果想要:

  • 在schema-on-read的情況底下,當資料寫入時也能夠驗證資料格式是否正確,那怎麼辦?
  • 在schema-on-read的情況底下,當資料讀取時也能夠驗證資料格式是否正確,那怎麼辦?

在schema-on-read情況之下想要驗證事件資料的正確性,基本上有兩種做法,Strong SchemaWeak Schema,前者是Apache Kafka與Apache Pulsar這類event broker所採用的方式,後者是Gregory Young在《Versioning in an Event Sourced System》所建議的方式。

***

Strong Schema

如圖2所示,在schema-on-read情況之下想要驗證寫入事件資料的正確性,有一種最簡單的做法就是資料寫入的時候在payload(event data)之前儲存事件的schema definition,如此一來寫入事件的程式(Serializer)就可以自動驗證事件的格式是否正確。此外,讀取事件的程式(De-serializer)也可以依據這個schema來讀取資料。

 



▲圖2:每一個事件夾帶一個schema definition,寫入時用來驗證事件格式的正確性,讀取時做為解析event data的依據。

 

圖2的作法雖然簡單,但是每一個事件都夾帶一筆schema definition,這些重複的schema definition在網路傳輸時會占用頻寬,儲存時會占用空間。所以衍生了「何不將這些schema definition集中起來放在一個地方儲存,serializer與de-serializer有需時再來這裡讀取就好?」的想法,這就誕生了Schema Registry這種軟體。

 

***

圖3是Apache Pulsar Schema Registry的運作示意圖,Schema Registry針對每一個Topic,保存曾經出現在裡面的事件版本及其schema。當Writer要寫入事件時,透過Pulsar client程式把事件的schema(SchemaInfo物件)向Schema Registry詢問該schema是否註冊過,如果沒有就註冊,如果有就得到一個代表該schema的serializer,然後透過它寫入事件。讀取事件時,同樣提供Schema Registry一個SchemaInfo物件,然後回傳一個de-serializer並透過它讀取事件。

基本上Schema Registry就是一個集中式的查表服務,採用Schema Registry之後,寫入資料庫中的事件就不需要像圖2一樣,每一筆事件都夾帶一個schema definition。只要把事件版本與schema definition向Schema Registry註冊,然後寫入事件時Schema Registry回傳的serializer會知道事件的版本編號,最後寫入資料庫只需要加上這個版本編號就可以。如圖3中事件前方綠色方框裡面的數字就代表該事件的版本編號


▲圖3:Apache Pulsar Schema Registry運作示意圖,參考《Apache Pulsar in Action

 

Schema Registry詳細運作還有很多細節,有興趣的鄉民可以參考Apache Pulsar或是Apache Kafka,這兩者都使用Schema Registry,只不過預設使用的Schema Registry軟體不同。在《Apache Pulsar in Action》書中針對Pulsar的Schema Registry運作原理有很詳盡的說明,有興趣的鄉民可以參考。

 

***

Weak Schema

Strong Schema在schema-on-read的情況下提供了寫入與讀取自動驗證事件格式的機制,感覺起來很不錯,可以減少人為錯誤。但Gregory Young在《Versioning in an Event Sourced System》書中卻建議採用另一種稱為Weak Schema的方式,其特性如圖4所示。

Weak Schema不使用Strong Schema的deserialization方式讀取資料,而是採用mapping的方式。什麼叫做mapping的方式?有點類似你把領域物件轉成DTO(data transfer object)往前端傳遞的時候,手動撰寫Mapper程式把Domain Object 轉成DTO的方式。只不過Weak Schema說的mapping,是把事件從資料庫映射(map)成領域事件物件的過程。

Mapping原則很簡單,只有三條,如圖4中所說明。

 


   ▲圖4:Weak Schema特性說明

 

採用Weak Schema的方式,事件本身就不需要加上版本號碼,因為只要依循Weak Schema的規則,不管資料庫中儲存多少事件的版本,讀取端一定能夠用mapping的方式把事件讀出來。此外。這種方式也不需要一個集中式的Schema Registry。

但是Weak Schema有一個很重要的限制,就是事件的欄位名稱不能改名(不能rename),還有就是如果有些必填欄位在某些事件版本被取消或是忘了填,則mapper程式要做特別的檢查,如圖4右下角程式所示。

 

關於事件版本控制的問題還有很多細節與相對應的方法,Teddy建議有興趣的鄉民可以仔細閱讀Gregory Young的《Versioning in an Event Sourced System》(這是一本電子書)。

 

***

哪一種比較好?

Strong Schema和Weak Schema哪種比較好?說實話Teddy也不知道。ezKanbna目前是採用偏向Weak Schema的方式,而且據Teddy的了解,EventStoreDB好像也沒有支援Schema Registry。EventStoreDB最早是由Gregory Young所開發,既然「祖師爺」建議採用Weak Schema,EventStoreDB不支援Schema Registry也是很合理的。

但如果今天的應用場景改成跨微服務訊息溝通,而不是單純的Event Sourcing,通常會使用Kafka或是Pulsar這種event broker軟體。在這種情況下使用Schema Registry(Strong Schema)就非常普遍。

但是,依據Gregory Young的說法,即使是在跨微服務訊息溝通的情況下,還是可以採用Weak Schema。就Teddy所知,Apache Pulsar的Topic在寫入與讀取的時候是可以設定不要檢查schema,這種方式就可以支援Weak Schema採用mapping的方式讀取事件。

兩種方式各有利弊,要用哪種方式就留給鄉民們自行決定。

***

 

下集預告

下一集談和事件版本管理相關的議題:行為版本管理

***

友藏內心獨白:搞了好一陣子才把這個議題弄清楚。

2022年7月9日 星期六

事件溯源(12):建立快照以加速聚合讀取

July 04 21:23~23:39

▲圖1:使用快照加速Aggregate讀取速度

 

前言

Event Sourcing系統在寫入端非常簡單且快速,但讀取因為需要從頭到尾逐一套用Aggregate instance所屬的所有領域事件以獲得最新狀態,經常讓人有一種「很慢」的感覺(錯覺?)。因此談到Event Sourcing除了套用CQRS將寫入與讀取模型分離以加快讀取速度以外,另一種常見加速的方式就是幫Aggregate的狀態建立快照(Snapshot)這一集就談這個議題。

 

***

快照原理

請參考圖1,Account Aggregate的原始資料儲存在Event Stream裡面,每次AccountRepository載入Account,必須從E1到EN傳給Account重新apply這些事件以計算出最新狀態。為了加速載入Account,幫Account的某個版本產生一筆快照,並將這個快照儲存至Snapshot Stream。下次AccountRepository再載入該Account時,會先到Snapshot Stream找出最新一筆快照,然後將快照的資料寫回Account,接著再從原本的Event Stream讀取快照之後所產生的新事件,然後將這些事件傳給Account再apply一次即可。

在圖1中,Snapshot Stream中有兩筆快照,第一筆 version=10,代表它是前10個事件的快照。第二筆快照是目前最新的快照,version=20,代表它是前20個事件的快照。當AccountRepository載入該Account時,它讀取Snapshot Stream最後一筆資料獲得最新的快照,因為快照的版本是20,所以AccountRepository接著從Event Stream的第21個位置開始讀取事件直到結束,然後把這些事件重新apply一次。在這個例子中,原本需要套用 N 個事件才可以獲得最新狀態,有了快照之後只需要套用 (N – 20) 個事件即可。

 

***

實作

實作快照的方法很多,一種常見的做法是讓Aggregate套用Memento設計模式產生快照,並自行從快照回復狀態。圖2是Memento設計模式類別圖,這個設計模式有三個角色:

  • Originator:需要保存快照資料的物件,以等一下要介紹的ezKanban例子而言,就是Tag Aggregate。
  • Memento:快照,在Memento設計模式裡面快照稱為Memento(備忘錄)
  • Caretaker:負則儲存快照的物件,ezKanban例子而言可以直接修改TagRepository程式碼讓它扮演Caretaker。如果想要遵守開放封閉原則不想修改原本已經可以動的TagRepository,則可以套用Decorator設計模式在Repository身上「外掛」快照的功能。

 

▲圖2:Memento設計模式類別圖

 

***

圖3是ezKanban中所定義的快照介面,它只有getSnapshot()與setSnapshot()兩個方法,前者用來產生快照,後者用來從快照回復狀態。


▲圖3:Memento介面

 

圖4為Tag Aggregate實作Memento介面的程式碼,第14行的TagSnapshot record就是快照本人,它代表Tag的目前狀態。第25行getSnapshot方法直接回傳一個新的TagSnapshot,第30行的setSnapshot方法則是將傳入的TagSnapshot身上的值寫入Tag,等於回復Tag的狀態。第18行是一個static factory method,可以直接從TagSnapsht得到一個Tag物件。

 

▲圖4:Tag Aggregate實作Memento介面

 

圖5為支援快照版本的TagRepository的save方法,第72行儲存原本Tag的領域事件,第74行~第84行則是判斷是否需要儲存快照。這個版本的TagRepository會依據使用者所設定的snapshotIncrement數值決定多少筆領域事件紀錄一次快照,如果snapshotIncrement等於100,則超過100筆會記錄一次快照。

 

▲圖5:TagRepository的save方法(支援快照版本)

 

圖6為支援快照版本的TagRepository的findById方法,第46行從Snapshot Stream載入舊的Snapshot領域事件,如果Snapshot不存在則第49行直接用原有的方式載入Tag;如果存在,第53行由Snapshot產生Tag,然後再載入Snapshot版本號加1的Tag原有領域事件(第54~55行),然後逐一apply它們(第56行)。

 

▲圖6:TagRepository的findById方法(支援快照版本)

 

以上產生與載入快照的方式,除了可以用在Aggregate身上,也可以用在Event Sourced Read Model上面,請參考上一集<事件溯源(11):撰寫JavaScript在EventStoreDB中產生自訂投影>所提到的加速EventStoreDB採用使用者自訂Projection做為Read Model的做法。

 

***

你真的需要快照嗎?

如果採用上述方式實作快照,因為套用Memento設計模式,所以Aggregate Root需要實作Memento介面。也就是說,從Clean Architecture的角度來看,不僅僅是位於Adapter層的Repository實作需要修改,連Entity層的Aggregate也要改。這雖然不是什麼大不了的修改,但畢竟讓系統變得更複雜。除非必要,否則不需要建立快照。

什麼叫做必要?這要從Aggregate的生命周期長短來判斷。以ezKanban為例,Card Aggregate代表看板上面的一項工作,它從出生到死亡(工作完成並歸檔),身上的領域事件可能大不了幾十個。在這種情況下,幫Card建立快照其實是不需要的。但是,如果是銀行的Account物件,代表使用者的交易紀錄。通常銀行客戶可能存在好幾年甚至數十年,累積的交易紀錄非常可觀。此種情況下,如果不透過快照來加速,可能會有效能的問題。

但是Event Sourcing系統中也不是只有快照和CQRS這兩種方式可以增加讀取速度,還有一種「年度結算(固定時間結算)」的做法和快照類似。以銀行為例,每一年度開始可能會把去年整年每一筆Account交易紀錄「濃縮(壓縮)」成一筆新的領域事件,並且從原本的event stream中將舊有的事件移至另一個event stream(或是每一年度產生一個新的event stream亦可)。

Teddy曾經在某本書(還是文章,忘了來源)看到一種說法:「領域事件少於一萬都不需要做快照。」這個數據鄉民們可以參考看看。

***

 

下集預告

下一集談另一個比較進階但卻很重要的議題:Event Versioning(事件版本異動)

***

友藏內心獨白:Memento設計模式終於派上用場了。

2022年7月8日 星期五

事件溯源(11):撰寫JavaScript在EventStoreDB中產生自訂投影

July 04 18:38~19:46

▲圖1:EventStoreDB的使用者自訂投影程式

 

前言

上一集介紹ezKanbana如何實作Projector以便在PostgreSQL資料庫中產生Read Model所需的資料,這一集回頭使用EventStoreDB的客製化Projection功能,建立另一種形式的Read Model。

 

***

透過JavaScript程式產生事件投影

Teddy在<讀取事件溯源(6):透過Projection查詢Event Store>介紹過EventStoreDB三種內建的Projection:$all、$ce-以及$et-。除了內建的Projection以外,EventStoreDB也支援使用者撰寫JavaScript程式來產生Projection。什麼,JavaScript程式?!沒錯,Teddy沒寫錯,鄉民們也沒看錯,就是JavaScript程式。EventStoreDB可以執行JavaScript程式,投影出使用者想要的event stream,再利用這個event stream產生讀取資料。

依據EventStoreDB官方文件的講法,它是一個為了寫入所開發的事件溯源資料庫,也是一個可以做為讀取模型的資料庫。EventStoreDB就是透過自訂Projection使得它可以做為取模型的資料庫。

 

圖1為ezKanban用來投影所有Board事件的JavaScript程式,第1行~第5行的fromStreams函數表示將這些event stream當成輸入來源。第11行~第13行則把所有接受到的事件,投影到另一個ReplayEvents-By-Board-[board id]的event stream裡面。這個使用者自行定義的event stream,可以做為ezKanban的replay功能的讀取模型。也可以拿這個event stream當成資料來源,傳給上一集<事件溯源(10):實作Projector>裡面的NotifyBoardContent程式,讓它即時投影出GetBoardContent所需的讀取模型。

 

***

 

採用EventStoreDB當成讀取資料庫的架構如圖2所示,以GetBoardContent查詢為例,它會直接從ReplayEvents-By-Board-[board id] event stream裡面讀出屬於某一個Board的所有事件,然後再呼叫NotifyBoardContent讓它即時重播這些事件並投影出目前狀態。


▲圖2:將EventStoreDB當成Read Database

 

ezKanban的GetBoardContent查詢並沒有使用這種方式當作讀取資料來源,但是團隊有用測試案例來驗證這種作法的效率,如圖3所示。

  ▲圖3:使用測試案例模擬將EventStoreDB做為GetBoardContent的讀取模型資料庫

 

***

效能比較

ezKanban團隊成員杜奕萱在她的碩士論文《套用命令與查詢責任分離以簡化聚合依賴:以 ezKanban 為例》比較使用PostgreSQL資料庫以及在上一集<事件溯源(10):實作Projector>產生JSON讀取模型的效能比較,如圖4所示。可以看出來,當卡片數量越多,套用CQRS的效益越大。

 

▲圖4:杜奕萱碩士論文所做的ezKanban在套用CQRS前後,GetBoardContent查詢的效能比較

 

在杜奕萱的碩士論文中並沒有比較這一集所介紹的採用EventStoreDB建立GetBoardContent讀取模型的方法,但ezKanban團隊在前幾天以圖3測試案例的形式做了非正式的比較,其結果如下:

  • 在8000張卡片的情況之下,花費的時間大約2179 ms。這個速度還是比沒有建立讀取資料要來得快很多(原本需要16841 ms)。
  • 團隊進一步模擬以1000筆領域事件建一次Snapshot為例,則產生GetBoardContent的時間可以縮短成87ms,非常接近在PostgreSQL資料庫中用JSON格式產生讀取模型的73ms時間。

***

EventStoreDB自訂投影的好處

若採用上一集<事件溯源(10):實作Projector>所介紹的架構來套用CQRS,則讀取端的Projector需要具備idempotent的特性,而且若讀取資料庫損毀,也需要想辦法讀取原本的領域事件再次投影出讀取模型(replay)。使用EventStoreDB的自訂投影,只要撰寫好JavaScript程式,EventStoreDB會負責產生投影,使用者不用煩惱idempotent的問題。至於replay也很簡單,因為所有的領域事件都保存在EventStoreDB所投影出的stream裡面,所以只要重頭讀取一次該stream的所有事件再逐一apply即可。雖然這種方式的效能沒有直接用JSON儲存View Model來得快,但如果有需要只要加上Snapshot,便可以大幅改善EventStoreDB做為讀取資料庫的效能。

至於CQRS的讀取資料庫要採用哪種方式,就要看鄉民們各自專案需求與技術能力而定。

***

下集預告

Event Sourcing與CQRS的核心觀念與技術介紹的差不多了,下一集之後談幾個比較進階一點的議題,先從幫Aggregate或是Read Model建立Snapshot(快照)談起。

***

友藏內心獨白:Write DB也可以是Read DB。

2022年7月7日 星期四

事件溯源(10):實作Projector

July 03 09:38~12:01


▲圖1:ezKanban套用CQRS之後的架構圖,箭頭方向代表data flow

 

前言

上一集介紹在ezKanban中套用CQRS簡化領域模型的例子,這一集將說明ezKanban如何實作Projector以便在PostgreSQL資料庫中產生Read Model所需的資料。

 

***

那些查詢需要快取?

讀取資料庫中專門為特定查詢所準備的資料,其實就是一種快取(Cache)。這些快取資料由Projector所產生,因此在資料庫層套用CQRS,必須要先決定:「需要幫那些查詢產生Read Model所需的快取資料?」

以ezKanbna為例,有兩個主要且比較複雜的查詢畫面,分別是圖2的GetDashboard,以及圖3的GetBoardContent。前者是使用者登入ezKanban之後看到屬於他個人的所有Team(畫面最左方)、每一個Team裡面的Project,以及每個專案中有那些Board的畫面。後者是使用者進入某一個Board之後,所看到屬於該Board的所有Workflow、Card與Tag。

由於GetBoardContent需要滿足多人線上協同合作的需求,讀取頻率相對來說很高,需要考慮讀取效率的問題。因此ezKanban幫GetBoardContent在讀取資料庫中產生快取資料,以優化效能並達到簡化領域模型(寫入模型)的目的。

至於GetDashboard是使用者個人看到的畫面,沒有多人線上協同合作的需求,比較沒有讀取效率的問題。因此ezKanban目前並沒有幫GetDashboard在讀取資料庫中產生快取資料,它的資料是即時從寫入資料庫中所產生。

從軟體架構的角度來看,Teddy之前提過可以在不同架構階層套用CQRS,例如API層(Adapter層)、使用案例層、領域模型層與資料庫層。以上述GetBoardContent與GetDashboard為例,GetBoardContent的CQRS套用到資料庫層,而GetDashboard只套用到使用案例層(GetDashboard本身是一個查詢使用案例)

在這裡Teddy提醒一點套用CQRS很容易誤會的點,就是在資料庫層套用CQRS是一個逐案探討(case by case)的情境,也就是說不是所有的查詢都需要在資料庫中準備一份快取資料。因為幫特定查詢準備快取資料這件事,需要撰寫特殊的Projector程式,而這件事也是一個開發成本。所以只有針對特別講求效能的查詢,或是為了簡化寫入模型這兩個目的,才需要在資料庫中產生快取資料(在ezKanban中目前只遇到這兩種情況)。

    

 ▲圖2:ezKanban的GetDashboard畫面

 

     

  ▲圖3:ezKanban的GetBoardContent畫面

 

***

NotifyBoardContent Projector

在ezKanban中產生GetBoardContent所需讀取快取的Projector稱為NotifyBoardContent,程式碼如圖4所示。由於NotifyBoardContent負責投影出整個Board裡面的所有資料,因此它需要去聽Board、Workflow、Card、Tag這四個Aggregate的領域事件,加起來一共有34個。

收到這些領域事件之後,NotifyBoardContent透過BoardContentStateRepository從資料庫中讀出原本的快取資料,然後更新這份快取資料,最後把快取資料寫回資料庫。參考圖4第60行~67行,當NotifyBoardContent收到WorkflowCreated領域事件之後,它先產生一個workflwoState物件代表這個新增的Workflow(第61行),然後以領域物件做為參數,重新apply一次這個事件,以便設定workflwoState的值(第62行)。

接著從資料庫中讀出BoardContentState(GetDashboard所需的快取資料,如圖5所示),將workflwoState加入BoardContentState,然後把BoardContentState回存到資料庫中,更新快取資料。如此便完成一次投影讀取資料的操作。

 


▲圖4:產生GetDashboard所需讀取快取的NotifyBoardContent Projector程式

 

圖5中的BoardContentState,不屬於DDD裡面的領域模型物件(不是Aggregate、不是Entity、不是Value Object也不是Domain Service),它是放在Use Case層的View Model物件,專門服務某特定畫面所需的資料結構。從程式碼中可以看出,第17行BoardState主要紀錄原本放在Board Aggregate的資料,第18行List<WorkflowState>記錄Board身上所有Workflow的資料,以及這些Workflow的順序。第19行儲存這個Board裡面所有Tag資料,第20行紀錄每一個Lane身上的所有Card資料。

由於BoardContentState是一個讀取模型,它的資料結構包含一整坨這個畫面所需要的全部資料都放在一起,因此讀取資料時只要下一個查詢條件就可以直接從資料庫讀出,加快讀取速度。

 

▲圖5:BoardContentStatet介面

 

BoardContentState轉成JSON(格式如圖6所示)存在PostgreSQL資料庫表格的jsonb欄位中,GetBoardContent讀出後直接傳給前端的React程式,React收到BoardContentState之後把它存在Redux並以此做為前端顯示資料的狀態。

 

▲圖6:BoardContentStatet的JSON檔案內容

***

Projector好難寫

從圖4可以看出NotifyBoardContent為了從34個領域物件投影出讀取模型,它的程式碼還挺複雜的,而且很多用來維持讀取模型狀態的程式碼和寫入模型中Board、Workflow、Card、Tag 這些Aggregate維持自身狀態的程式碼幾乎相同,所以會有重複程式碼壞味道的問題產生。

在ezKanbna中因為剛好套用DCI(Data Context Interaction)架構,因此NotifyBoardContent可以重複使用寫入模型中用來更新Aggregate狀態的函數,解決重複程式碼的問題。但是在一般軟體開發專案中,如何設計與撰寫Projector的確是一個需要注意的問題。

 

***

下集預告

除了本集所介紹的NotifyBoardContent這種接收寫入模型的領域事件然後在資料庫中產生Read Model的方式以外,還有其它不同的做法。例如,使用關聯式資料庫建立Indexed View就是非常簡單且常用的產生讀取模型作法。另外像是Database Replication(利用資料庫內建的複製功能,將主要資料庫複製到次要資料庫)也是一種分散讀取負載的做法(把次要資料庫當成Read Database使用)。還有就是使用Change Data Capture (CDC)工具,自動擷取出資料庫中異動的資料(不依靠領域事件產生資料異動,而是直接在資料庫層級攔截資料異動紀錄),然後再投影出讀模型。

以上三種方式屬於傳統關聯式資料庫所支援的產生讀取模型方法,和Event Sourcing比較無關。在下一集中,Teddy將介紹EventStoreDB所支援的另一種直接在事件溯源資料庫中產生讀取模型的方法。

***

友藏內心獨白:雖然套用了DCI,ezKanban的NotifyBoardContent也是寫了好幾天。


2022年7月6日 星期三

事件溯源(9):套用CQRS簡化領域模型

July 02 18:35~19:12;23:15~24:00;July 03 00:00~13:07

▲圖1:ezKanban Core Domain Model(簡化版)

 

前言

上一集提到CQRS可以簡化設計,這一集以ezKanban core domain為例,說明套用CQRS之後如何簡化原本領域模型之間Aggregate(聚合)的關係。

 

***

雙向關聯造成不必要的複雜度

圖1是ezKanban套用CQRS之前Core Domain的領域模型簡化版,一共有四個Aggregate:Board、Workflow、Card、Tag。其中有兩對Aggregate保持雙向關聯,分別是:

  • Board與Workflow:一個Board可以有多個Workflow,而且必須記錄每一個Workflow在Board上面的順序(order)。此外,Workflow身上紀錄boardId,讓它知道自己屬於哪一個Board。為了記錄Board身上每一個Workflow的順序關係,Board身上有List<CommittedWorkflow>屬性,其中CommittedWorkflow是一個association class,身上有boardId, workflowId, order這三個屬性。
  • Workflow的Lane與Card:一個Lane上面有多張Card,而且必須記錄每一張Card的順序(order)。此外,Card身上也紀錄著workflowId與laneId,讓它知道自己屬於哪一個Workflow的哪個Lane。為了記錄Lane身上每一張Card的順序,Lane身上有List<CommittedCard>屬性,CommittedCard也是一個association class,身上有cardId, laneId, order這三個屬性。

 

接下來以Board與Workflow的關係為例,說明這個雙向關聯對於領域模型造成什麼影響。請參考圖2,為了維持這個雙向關聯,Workflow與Board之間需要狀態同步,CreateWorkflow之後,需要通知Board在它身上加入這個Workflow。在DDD中,Aggregate之間的狀態同步是「狀態最終一致性」。換句話說,為了維持這個雙向關聯,領域模型的實作變得比較複雜(需要維持狀態最終一致性)。

 

▲圖2:為了維持雙向關聯Workflow與Board必須達成狀態最終一致性

 

仔細想一想,為什麼Workflow與Board之間需要維持雙向關聯?因為Board需要知道它身上有多少個Workflow,以及這些Workflow的順序。繼續追問下去,那麼為什麼Board需要知道它身上的Workflow與順序?是Board的業務邏輯需要這些資訊嗎?完全沒有。這些資料是為了顯示用途而存在。如圖3所示,ezKanban的GetBoardContent顯示Board裡面有3個Workflow。也就是說,為了顯示(查詢)用途而導致領域模型增加不必要的複雜度。

 

▲圖3:包含三個Workflow的Board畫面

 

***

Eric Evans怎麼說

在《Domain-Driven Design: Tackling Complexity in the Heart of Software》書中作者Eric Evans提到:

It is important to constrain relationships as much as possible. A bidirectional association means that both objects can be understood only together. When application requirements do not call for traversal in both directions, adding a traversal direction reduces interdependence and simplifies the design. Understanding the domain may reveal a natural directional bias. (盡可能地限制關係很重要。雙向關聯意味著兩個物件只能一起理解。當應用程式不需要雙向遍歷時,採用單向遍歷可以減少相互依賴並簡化設計。了解(問題)領域可能會揭示一種自然的方向偏差。)

如果在問題領域中單向依賴就可以解決問題,在領域模型中就不需要維持雙向依賴。一開始ezKanban並沒有套用CQRS,所以它的領域模型很自然地需要同時滿足寫入(Workflow身上有boardId)與讀取(Board身上有List<CommittedWorkflow>屬性)的需求。在傳統物件導向分析與設計(OOAD)中,因為沒有Aggregate的觀念,所以這種雙向關係並不會造成什麼大問題。因為Board與Workflow直接可以透過記憶體參考而存取對方,因此Workflow狀態改變時Board立即可得知,不需要透過領域事件做到狀態最終一致性。但是在DDD中,因為ezKanban把Board與Workflwo設計成兩個不同的Aggregate,所以這種混合寫入與讀取的單一領域模型,從寫入的角度來看,便產生不必要(透過領域事件達到最終一致性)的複雜度。

Eric Evans在書中提到一些簡化關聯的做法,但是並沒有從讀寫分離的角度來探討如何簡化領域模型的關聯。

***

套用CQRS簡化模型複雜度

Teddy剛剛分析過,Board身上的List<CommittedWorkflow>是為了查詢而存在,寫入模型並不需要這個資料結構。套用CQRS之後,圖1中ezKanban領域模型的兩個為了讀取模型而存在的association class(CommittedWorkflow與CommittedCard)就可以直接拿掉,如圖4所示。

簡化後的寫入模型,省去了不必要的關聯,也去除了不必要的狀態最終一致性。

 


▲圖4:ezKanban套用CQRS之後的寫入模型

 

但是問題來了,ezKanban還是需要知道Board身上有多少個Workflow,以及這些Workflwo的順序。在寫入模型中拿掉List<CommittedWorkflwo>之後,這個資料要從哪裡來?這就要靠CQRS的Query Model來記錄這個關聯性,如圖5所示。

 

▲圖5:ezKanban套用CQRS之後的讀取模型

 

現在剩下最後一個問題:「怎麼產生讀取模型所需的資料?」請參考圖6,在讀取模型中必須撰寫一支用來在讀取資料庫中產生Read Model所需資料的Projector程式。在它會監聽Write Model所發出的領域事件,然後依據這些領域事件在Read Database中投影出Read Model所需的資料。這種資料被稱為「非正規化」或「物質化」資料,為了快速讀取可以允許重複的資料存在。

以ezKanban為例,GetBoardContent查詢所需的資料是一個代表整個Board所有資料的JSON物件,存在PostgreSQL資料庫的jsonb欄位。GetBoardContent查詢只需要下一個SQL指令就可把整個Read Model所需的資料從PostgreSQL讀出來,不需要下任何的join條件,所以查詢的速度很快。

 

▲圖6:ezKanban套用CQRS之後的架構圖,箭頭方向代表data flow

 

在這裡有兩個重點要注意,首先Write Database與Read Database之間的狀態是最終一致性,也就是說Read Model不一定會有最新的資料,這一點在系統設計時被需要考量進去,否則可能會造成使用者體驗不佳。例如,使用者剛剛才下一筆訂單(存在於Write Database中),但在訂單查詢畫面(從Read Database讀取)中卻查不到這筆訂單的資料。所以Teddy常說CQRS雖然簡化寫入與讀取模型內部的複雜度,但卻把複雜度轉換成兩個模型之間狀態同步的問題。至於如何取捨,就要看實際的業務需求與應用情境而定。

第二個重點是,這支Projector程式雖然是屬於Read Model,但它做的工作是「產生Read Model」,也就是說它是負責寫入Read Model的人。它的程式邏輯,可能有部分,甚至很多,和Write Model的Aggregate身上的邏輯互相重複。另外,因為它可能收到重複的領域事件(在分散式系統中,事件傳遞通常只能滿足at least once,不容易做到exactly once),因此它需要滿足idempotent,否則可能投影出錯誤的Read Model。最後,因為Read Model可能因為某些原因導致本身的狀態錯誤或是被刪除,因此Projector必須有能力能夠從頭重建Read Model,通常是藉由replay所有相關的領域事件來達到此功能。

 

***

下集預告

看完CQRS的概念說明,下一集將介紹ezKanbana如何實作Projector,在PostgreSQL資料庫儲存Read Model。

***

友藏內心獨白:頭快爆炸了嗎XD。