l
顯示具有 Programming 標籤的文章。 顯示所有文章
顯示具有 Programming 標籤的文章。 顯示所有文章

2023年1月4日 星期三

也許我不會刪除你:重複程式碼不等於壞味道

January 04 00:00~01:10


前言

今年一月上班沒幾天就要過農曆年,去年底在安排泰迪軟體2023年課程表的時候,索性就把整個一月空下來。本來想要找機會再溜出國玩,後來因為去年年底太忙,沒時間安排行程,乾脆把時間拿來在家裡讀書算了。

這幾天同時間讀了幾本書,其中有一本《Five Lines of Code: How and when to refactor》滿有意思的。這是一本討論Refactoring(重構)的書,有別於傳統的重構,書中採用Rule-Based Refactoring,藉由提供幾條非常具體的規則,讓開發人員知道何時以及如何重構。例如書中的第一條規則就是這本書的書名:每一個method的程式碼不可以超過5行(不包含括號)。

今天Teddy要談這本書提到的另一個重構中經常遇到的問題,就是Duplicated Code(重複程式碼)。傳統上認為重複程式碼是Bad Smell(壞味道、怪味道),應該要將其斬草除根。關於這一點,基本上沒什麼爭議。如果程式中有很多重複程式碼,一旦需求改變涉及到這些程式碼,開發人員就需要修改每一處的重複程式碼,否則系統就會發生錯誤。這就造成另一個壞味道:Shotgun Surgery

***

再探重複程式碼

這幾年Teddy因為學習Clean Architecture與Microservices,發現透過重複程式碼來避免模組之間的依賴,反倒是一種很常見的方法。例如,Clean Architecture的跨層原則要求Entities Layer的物件離開Use Cases Layer時必須轉成另外一種資料結構;若是往前端傳遞則轉成DTO(Data Transfer Object),往資料庫傳遞則轉成PO(Persistent Object)。雖然DTO與PO與其所相對應的Entities Layer物件並不一模一樣,重複的部分只有資料,但廣義來看這也可以視為一種重複程式碼。


Five Lines of Code: How and when to refactor》書中對這一點解釋的很好:

Sharing code increases global behavior-change velocity, while duplicating code increases local behavior-change velocity.

針對某個功能,如果全部程式碼都共用(沒有重複程式碼),那麼當系統行為改變的時候,只要改一個地方即可,因此「global behavior-change(全域行為改變)」的速度會很快。反之,如果同一個功能每次用到它的時候都將其複製一份,那麼這份複製出來的程式碼就與原本的「本尊」獨立,去除了耦合。在這種情況下,「local behavior-change(區域行為改變)」的速度就會很快。


從生物學的角度來看,獨立的區域,例如島嶼、沙漠或高山,經常會演化出「特有種」生物。一開始這些生物的起源是相同的,但是為了應付「區域性需求」,逐漸演化出不同的特徵。 程式碼也經常如此,因為被共用的程式碼在使用它的地方可能同時存在不同的區域性需求(書中稱為local invariants),此時如果堅持「共用」,很可能導致需要回頭修改共用程式碼,讓它具備適應不同區域性需求的能力(透過設定或依賴注入),因而增加共用程式碼的複雜度,最後可能會複雜到降低它的可讀性,因而反倒造成可修改性下降(原本共用希望提升global behavior-change速度,但在這種情況下,可能反倒降低global behavior-change速度)。

***

Shotgun Surgery的問題勒?

討論至此,Teddy也不是鼓吹盲目使用「複製、貼上」來去除模組之間的依賴這樣就好棒棒,「利用重複性去除依賴」這個方法,還是要從軟體架構邊界這兩個角度來考慮。假設你發生重複性的邊界是在同一個method裡面,這種重複性幾乎可以斷定是壞味道,可以用Extract Method將其移除。如果發生重複性的邊界是同一個package,則也有很大的機會是壞味道。但是如Teddy前述的情況,在Clean Architecture的跨層原則之下,重複性的邊界已經是架構階層之間,此時使用重複性去除依賴爭議就比較小。

另一個常見的情況是在微服務架構中,下游微服務透過聽取上游微服務的事件,在本地端建立讀取模型(Read Model)以隔離兩個微服務之間的執行期間依賴。

 

***

結論

Teddy在開發ezKanban的過程中也遇到很多「是否要共用,還是利用重複性去除依賴」的設計決定,例如ezKanban支援Event Sourcing與State Sourcing這兩種狀態儲存方式,一開始ezKanban的報表先支援State Sourcing,然後再支援Event Sourcing。為了支援這兩種儲存方式,團隊在Repository的實作中套用Pluggable Adapter設計模式,整個設計簡單易懂。

等開發完成之後,團隊發現有不少報表除了撈資料的方式不同(一個下SQL另一個操作event streams),其餘產生報表的計算邏輯大致上是相同的。為了去除這些重複性,又花了不少時間重構系統。重構完成後,去除重複程式碼,但付出的代價就是設計變得更間接(因為多了一層抽象介面),也沒那麼直覺。

Teddy覺得這是一個有趣的議題,重新省視「程式碼共用」這件事,共用並不一定都是好的,想想你的中台…XD。

 

***


友藏內心獨白:To share, or not to share, that is the question。

2022年8月17日 星期三

我可能不會用你,Event Sourcing + CQRS?!(下)

August 17 21:54~22:49

▲圖1:Repository介面

 

前言

在〈我可能不會用你,Event Sourcing + CQRS?!(上)〉Teddy談到Event Sourcing在寫入模型(Write Model/Command Side)並不會比State Sourcing要複雜,甚至更簡單。

但事情通常都是有一好沒兩好,寫入比較簡單,讀取就比較困難,今天討論Event Sourcing在讀取模型的情況。

***

跨Aggregate的資料查詢

在套用DDD(領域驅動設計)之後,Aggregate Root負責發出領域事件,然後透過Repository儲存它的狀態。當採用Event Sourcing時,一個aggregate instance在event store中會有一個event stream用來保存它的所有領域事件。每一個Repository(DDD裡面的Repository設計模式)負責儲存與載入「單一Aggregate」,其介面如圖1所示。

圖2為ezKanban系統core domain的領域模型,它包含Board、Workflow、Card、Tag四個Aggregate。現在問題來了:如何查詢一個Board裡面有多少個workflow?多少張卡片?多少種Tag?」

如果把DDD Repository設計模式定位為「專門在Write Model中用來存取單一Aggregate的介面」,而且Write Model不需要維持「為了查詢而存在的關聯」,那麼在Write Model中,Board並不知道它身上有多少個Workflow以及Tag,而Workflow也不知道它身上每一個Lane有多少張Card。

 

▲圖2:ezKanban領域模型

 

不用維持雙向關聯之後的Event Sourcing更進一步簡化寫入模型,但查詢就比較困難,特別是針對跨Aggregate之間的查詢。如果不管效率問題,開發人員還是可以透過EventStoreDB內建的Projection功能,讀出領域事件然後在記憶體中「拼出」read model,但這樣執行速度顯然會比較慢。所以在套用Event Sourcing的情況下,針對特別的查詢畫面或報表,通常會特別設計一個View Model以加速讀取速度。

在State Sourcing的情況下,如果是使用關連式資料庫,查詢可以直接下SQL找出所需要的資料。但這並不表示State Sourcing就不會有查詢效率的太慢的問題。同樣地,針對個別查詢畫面,State Sourcing也經常會建立Read Model來加速查詢速度。換句話說,額外建立Read Model並不是Event Sourcing的專利,在State Sourcing系統中,下SQL、Join、Create View都是一種產生Read Model的方法。

***

Event Sourcing到底

傳統上在Event Sourcing系統中套用CQRS,Read Model資料庫的選擇可以和Write Model相同或不同(好像是廢話)。例如Write Model如果用EventStoreDB,Read Model可能用關連式資料庫或NoSQL。但是有另一派的做法是:「既然要Event Sourcing,就Event Sourcing到底,Write Model與Read Model都Event Sourcing。」這是什麼意思?請參考〈事件溯源(11):撰寫JavaScript在EventStoreDB中產生自訂投影〉,Teddy介紹過EventStoreDB可以透過撰寫JavaScript程式讓資料庫自動產生給查詢使用的projection(自訂projection),然後用這些projection來產生Read Model。在這種情況下,Write Database與Read Database就可以合併成一個。

***

 

結論

看到這裡鄉民們可能會覺得:「Teddy你還是沒說Event Sourcing + CQRS是不是比State Sourcing要難啊?」還是那句老話,沒有比較難,只是不一樣。哪裡不一樣:

  • 乍看之下,Event Sourcing寫入比較簡單,讀取比較難,State Sourcing則相反。
  • 實際上則是「能量不滅」,各有各的困難之處。Event Sourcing不須要設計資料庫schema,感覺好棒棒。但設計資料庫schema的問題並沒有消失,只是轉成設計領域事件schema。

Teddy寫這兩篇文章的目的不是要倡導大家使用Event Sourcing,只是要提醒,它就是一種儲存狀態的方式,它沒有比State Sourcing難。如果你找到適合Event Sourcing的應用場景,使用它可以簡化系統設計。如果你的應用場景不需要記錄所有狀態異動,而且你也不熟悉Event Sourcing,那麼就用已經習慣的State Sourcing就好,不用趕流行。

但是,多了解一種狀態儲存方式也沒什麼壞處,難保哪一天合適的應用場景出現了,此時不知道Event Sourcing可能會設計出過於複雜的系統。

***

友藏內心獨白:學習新技術就是為了未來有能力可以看到forces。

2022年8月16日 星期二

我可能不會用你,Event Sourcing + CQRS?!(上)

August 16 10:03~11:33


▲圖1:Clean Architecture四層架構

 

前言

上周末上完【事件溯源與命令查詢責任分離架構實作班】首發團,有一位認識近十年的老學員告訴Teddy:「這是我上過泰迪軟體所有的課裡面最難的一門課。以前的課我回去都可以直接挑選某部分在工作上應用,這個Event Sourcing加上CQRS我一下子想不到可以用在哪裡。」

Event Sourcing真的比State Sourcing(OR-Mapping)要困難嗎?這次分兩集談一下這個問題,這一集先從寫入(Command/ Write Model)來比較兩者的異同,下一集再來討論讀取。

 

***

只是不一樣

Teddy覺得這位老學員之所以會覺得難,主要的原因是在課程中Teddy把近幾年所學關於Event Sourcing, CQRS, DDD, Clean Architecture的全部重點濃縮在兩天內講完,所以資訊密度比較高。至於Event Sourcing本身Teddy認為並沒有比State Sourcing要難,困難的地方在於你會想應用Event Sourcing的場景通常是在分散式系統或套用微服務架構,是這個「分散式環境」讓你誤以為Event Sourcing比較難,而非它本身真的比較難

和State Sourcing相比,Event Sourcing只是做法不同,並沒有比較難;它們只是一種儲存系統狀態的方法。

 

***

從Clean Architecture看Event Sourcing

請參考圖1的Clean Architecture架構圖:

Entity Layer:Entity Layer首先反映問題領域的重要概念與關聯,理論上這一層的物件根本不管物件如何儲存。但在程式實作面,Event Sourcing對Entity Layer的實作方式的確會產生影響,但也僅止於影響Aggregate Root,程式撰寫風格要改成如圖2的event sourced coding style:任何改變狀態的public method,都呼叫apply方法,並傳給它一個(或多個)代表狀態改變的領域事件。然後在event handler(圖2中的when方法)中實作程式邏輯。這種程式撰寫風格,也可以適用state sourcing,所以也不是說寫成這樣就一定要用event sourcing。

 

▲圖2:Event Sourced Aggregate Root撰寫風格

***

 

Use Case Layer:請參考圖3,Use Case透過注入的Repository來存取Aggregate,Use Case並不知道所注入的Repository是一個State Sourcing Repository還是Event Sourcing Repository。也就是說,是否採用Event Sourcing並不會影響到Use Case。


▲圖3:Use Case範例

***

Interface Adapter Layer:請參考圖4,Interface Adapter Layer左方是State Sourced Repository的實作,右方是Event Sourcing Repository實作。從寫入(Write Model)的角度來看,如果是採用OR-Mapping的State Sourced Repository,由於還要寫一堆ORM設定或是下SQL寫入資料庫,實作方式其實比Event Sourced Repository還要複雜,並沒有Event Sourcing的儲存方式比State Sourcing要困難的問題。

 

▲圖4:Interface Adapter Layer與DB & Driver Layer

 

***

DB & Driver:請參考圖4,DB & Driver Layer左方是State Sourced Database,右方是Event Sourced Database。如果左方採用關聯式資料庫,將Aggregate寫入資料庫表格的同時,需要在同一個交易中把領域事件也一併寫入資料庫的Outbox表格。至於右方的Event Store Database,由於只需要寫入領域事件,不需要像關連式資料庫一樣寫入時套用Transactional Outbox設計模式,所以反而比較簡單。

***

結論

經過以上分析,從讀寫分離的角度來看,Event Sourcing在寫入端其實比State Sourcing還要簡單。但就像Teddy經常說的「能量不滅」,寫入比較簡單,通常讀取就比較困難。下一集從讀取模型來比較Event Sourcing和State Sourcing。

***

 

工商服務

最後打個廣告,10月份【事件溯源與命令查詢責任分離架構實作班】招生中。

***

友藏內心獨白:Event Sourcing是一種第一次會痛,第二次會爽的技術 XDD。

2022年7月26日 星期二

事件溯源(19):在InMemoryRepository實做樂觀鎖

July 26 15:52~16:39

▲圖1:樂觀鎖測試案例

 

前言

這系列文章原本是Teddy為了製作【事件溯源與命令查詢責任分離架構實作班】課程範例而撰寫,課程範例已經完成,這系列文章也已經連載結束。前幾天Teddy回頭把課程範例改成新的寫法,修改之前先跑測試,居然有一個錯誤!

仔細一看才想起來Teddy之前練習的時候注入在測試案例中InMemoryTagRepository,它沒有支援樂觀鎖,所以原本的樂觀鎖定測試案例會失敗,只要換回正常的Repository就好了,這個問題以前也遇到過。

但是這次Teddy突然想到:「為什麼InMemoryRepository不能支援樂觀鎖?」花了幾分鐘改一下Code,測試案例就通過了。今天就追加一篇,談如何讓InMemoryRepository支援樂觀鎖。 

***

實作樂觀鎖

Teddy在<事件溯源(7):樂觀鎖>中介紹過如何在關聯式資料庫與事件溯源資料庫實作樂觀鎖,基本上就是要在Aggregate身上加一個Version欄位,每次儲存Aggregate的時候比對它身上Version的數值與資料庫中的數值是否相等。如果相等,就代表這個Aggregate上次從資料庫讀出之後並沒有其他人寫入,因此它目前的版本是最新的,可以直接儲存到資料庫中。反之,則代表目前Aggregate的版本比較舊,無法儲存,系統要丟出樂觀鎖定失敗例外。

請參考圖1測試案例,從tagRepository根據相同的tagId拿出tagV1與tagV2兩個相同的物件。先把tagV1改名後儲存起來,接著再儲存tagV2,此時tagV2身上的Version數值會小於tagRepository所儲存的數值,因此會丟出RepositorySaveException例外。

 

首先修改InMemoryTageRepository的findById方法,如圖2所示。原本InMemoryTagRepository將Tag儲存在List裡面,findById回傳的記憶體中Tag的參考(reference)。這種直接回傳記憶體參考物件無法測試樂觀鎖,因為圖1中tagV1和tagV2會參考到同一個tag,也就是說改了tagV1會同時改變tagV2的值。所以findById要改成回傳一個新的Tag物件,而不是原本Tag物件的參考。


▲圖2:修改InMemoryTagRepository的findById方法以支援樂觀鎖

  

其次修改InMemoryTagRepository的save方法,如圖3所示。如果要儲存的tag已經存在InMemoryTagRepository,而且它的版本不等於記憶體中的版本,則丟出RepositorySaveException。反之,先把tag從記憶體中移除(如果不移除,相同的Tag會出現在InMemoryTagRepository兩次),然後將它的版本加1,然後把它儲存起來,最後清掉tag身上的領域事件。


▲圖3:修改InMemoryTagRepository的sava方法以支援樂觀鎖

 

就這樣,這麼簡單。

 

***

結論

本集介紹如何讓InMemoryRepository也具備樂觀鎖,但在這裡Teddy實作的InMemoryTagRepository只支援State Sourcing的儲存方式,並沒有支援Event Sourcing。如果是要實作InMemoryEventSourcingRepository,基本上也不會太困難,應該只需要:

  • 把資料結構由List改成Map<String, List<DomainEvent>>,Map的Key是Event Stream Name,Value是Aggregate的領域是件。至於Aggregate的版本就是List<DomainEvent>的大小。
  • 在儲存Aggregate的時候,不需要更新版本號碼,因為讀取(fndById)的時候Aggregate的版本號碼就是它所屬的List<DomainEvent>的大小。

InMemoryEventSourcingRepository的實作就交給鄉民自行練習。

 

***

友藏內心獨白:這一集算番外篇。

2022年7月15日 星期五

事件溯源(18):實做Idempotent

July 9 14:20~16:20


      ▲NotifyBoard實做Idempotent架構圖

 

前言

Teddy在<事件溯源(16):分散式系統的事件語意與Idempotent>介紹過為什麼Event Handler需要具備Idempotent,這集以ezKanban系統中產生GetBoardContent的Event Handler—NotifyBoard為例(請參考<事件溯源(10):實作Projector>),介紹如何實做Idempotent。

 

***

記住你做過的事

實做Idempotent可以從兩個方向著手:

  • 操作本身即是Idempotent:如果一個系統的操作本身就是Idempotent,哪麼Event Handler就不需要特別處理看過的事件,只要確定事件順序正確,收到事件之後閉著眼睛執行一次即可。例如,delete操作本身是Idempotent,收到CardDeleted事件(卡片被刪除)直接套用一次即可。就算重複執行相同的CardDeleted也不會造成系統狀態錯誤。
  • 記憶已處理過的事件:在一般通用系統中,不太容易把所有系統操作都設計成具備Idempotent特性,因此Event Handler需要紀錄它曾經處理過的事件代號,然後每次收到事件之後要去查看該事件是否已經處理過了。如過是,則丟棄該事件;若否,才處理該事件然後把事件代號紀錄下來。這裡有兩個地方要注意,首先事件代號需要唯一,不可重複,才可判斷是否曾經處理過。其次,處理事件造成的系統狀態改變和儲存事件這兩件事,必須要在同一個交易(transaction)中完成,否則可能造成系統狀態改變但卻沒有把處理過的事件紀錄下來,這樣下次再收到相同事件變會重新執行一次,就沒有達到Idempotent。

在本文中Teddy要採用第二種方式實做Idempotent。

 

***

先看測試結果

鄉民們讀到「在同一個交易中儲存系統狀態改變與事件代號」這句話,是不是有種似曾相似的感覺?沒錯,這和Teddy在<事件溯源(4):將Aggregate儲存至Outbox Store>介紹過的方法是類似的。

圖1是NotifyBoardContent的測試案例,產生一個Board,一個Workflow,三個Stage,然後新增一張卡片。

 

▲圖1:NotifyBoardContent測試案例

 

圖1測試案例所投影出的BoardContentViewModel如圖2所示。

 

▲圖2:BoardContentViewModel(JSON檔案)

 

除了在資料庫投影出BoardContentViewModel,因為NotifyBoardContent支援Idempotent,所以資料庫的Idempotent表格同時也紀錄了NotifyBoardContent所處理過由測試案例所產生的七個事件,如圖3所示。這七個事件分別是:BoardCreated、BoardMemberAdded、WorkflowCreated、StageCreated、StageCreated、StageCreated、CardCreated。

 


▲圖3:資料庫Idempotent表格紀錄Event Handler讀過哪些事件


***

 

實作

NotifyBoardContent的project方法如圖4所示,57行縮起來的switch敘述是負責投影的程式邏輯,這邊要關注的是:

  • 第53~54行:呼叫boardContentStateRepository的isEventHandled方法判斷領域事件是否已經處理過,如果已經處理過就直接return。
  • 第241行:更新boardContentState身上的IdempotentData資料結構,它用來記錄Event Handler正在處理哪一個領域事件,如圖5所示。
  • 第242行:把boardContentState儲存到資料庫。boardContentStateRepository.save方法會在同一個交易中將boardContentState儲存在board_content表格中(圖2的那個JSON檔案),以及將IdempotentData儲存在idempotent表格中,如圖6所示。

 

▲圖4:NotifyBoardContent的project方法

 

 

▲圖5:IdempotentData類別

 

 

▲圖6:儲存boardContentViewData與IdempotentData要在同一個交易中完成

 

***


好像沒有很難?

看完上面實作方法,感覺要讓Event Handler達到Idempotent好像沒有很難。但是,還是有一些實做細節需要考慮。在ezKanban中,boardContentViewData與IdempotentData都被存放在PostgreSQL資料庫,因此可以用關聯式資料庫的transaction確保這兩個操作的狀態一致性。但如果鄉民將read model儲存在NoSQL資料庫,例如document-based NoSQL資料庫,這種資料庫不一定會提供「跨document」的交易控制功能,這時候就可能需要把處理過的事件編號儲存在代表read model狀態的document身上。

 

***

 

第一季結束

這一系列寫到這裡也就差不多了,Event Sourcing與CQRS的重點還有實做細節都交代過,之後如果還有想到什麼再補充說明。

 

***

友藏內心獨白:感謝收看。

2022年7月14日 星期四

事件溯源(17):讀取手刻Event Store所儲存的事件

July 7 18::09~19:21;21:15~23:23;July 8 13:45~16:27

▲在Event Store儲存Checkpoint

 

前言

雖然使用EventStoreDB這種專為Event Sourcing與CQRS所設計的資料庫可以減少許多開發工作,但實務上開發人員可能因為公司要求或專案限制,只能使用關聯式資料庫。在這種情況下,就必須要自己用關聯式資料庫模擬Event Store。

Teddy在<事件溯源(4):將Aggregate儲存至Outbox Store>介紹過如何使用關聯式資料庫同時儲存傳統ORM的表格資料與領域事件,但還沒談到如何讀取這些事件的方法,今天就來介紹這個議題。

 

***

保證事件的儲存順序

Teddy以ezKanbna使用的Message DB這個開源軟體(https://github.com/message-db/message-db)為例,介紹它如何在寫入時確保所有事件的順序。圖1是Message DB用來建立儲存事件表格的指令,第7行global_position欄位在每次新增一筆資料的時候帶入該欄位目前最大的數值加1(簡單想成這是一個自動增加的欄位),透過這個欄位來維持事件的順序。

沒了,就這麼簡單。

 

▲圖1:Message DB產生儲存事件表格的指令

 

***

至少一次(At least once)

Message DB本身只包含使用PostgreSQL當成Event Store所需的程式碼,並沒有提供客戶端程式,使用者必須要自行撰寫,也就沒有直接支援at least once。Message DB原本屬於Eventide Project裡面的一個模組(子專案),Eventide是支援Ruby語言的Event Sourcing與Pub/Sub開源軟體,其使用方法可參考Eventide官方文件

Eventide的Consumer程式是用Ruby開發,Teddy沒有用過,從它的官方文件(圖2)也看不出來是否有提供at least once的功能。圖2第4點提到:「Consumer會定期將客戶端讀取的位置自動寫入到後端」,在這種情況下,假設後端已讀取資料但尚未處理,但Consumer卻將客戶端讀取的資料位置寫入後端(代表資料被讀走),這樣可能會造成訊息遺失。

 

▲圖2:從Eventide的文件看不出來是否有支援at lease once

 

***

 

先不管Message DB的「原生家庭」Eventide提供的Ruby客戶端程式否有支援at least once,Teddy在此說明如何做到at least once的常見方法。在《Enterprise Integration Patterns》提到的方式是使用Transactional Client設計模式,它的概念就是在上一集<事件溯源(16):分散式系統的事件語意與Idempotent>中Teddy介紹的Pulsar做法,客戶端確定事件處理完畢之後要向伺服器發出ack,類似資料庫的commit指令,完成這筆交易,如圖3所示。

 

▲圖3:Consumer向Server發出ack之後被讀出的事件才會從Topic刪除

 

***

 

如果是自己實作讀取資料庫中事件表格的驅動程式,要怎麼做出類似效果?做法也不難,只要針對每一個Consumer儲存一個checkpoint代表它目前讀取到第幾個事件。當Consumer把事件讀走且處理完畢之後,再把這個checkpoint加1(如果是批次處理事件則可以一次加N),代表它已經讀走某個事件了。事件並沒有真的從資料庫中被刪除,而是用checkpoint的數值代表每一個Consumer讀取的位置。

基本上checkpoint就好像讀取陣列時所使用的index,一個event stream同時間可支援多個Consumer讀取,每一個Consumer讀取的進度(位置)都不相同。所以每個Consumer要取唯一的名字,用這個名字當作checkpoint的名字來記錄每個Consumer的讀取進度

採用這種作法,如果要重讀整個stream,只要把checkpoint刪掉就可以了。很簡單,對不對。

現在問題來了,這個checkpoint要存在哪裡?參考EventStoreDB的官方文件作法,可以把checkpoint存在server端或client端,形成兩種不同的subscription(Consumer):

  • Persistent Subscription:在Server端儲存checkpoint,如圖4所示第113行呼叫ack方法就可以標註哪一個事件已經處理完畢並更新Server端checkpoint的數值。但是因為EventStoreDB的Persistent Subscription支援上一集介紹過的Competing Consumer(競爭消費者)且有自動rerety(重送事件)的功能,所以並不保證事件的順序(Consumer收到的事件順序可能和資料庫中儲存的順序不一樣)。EventStoreDB的文件建議,如果要保證順序請使用它的Catch-up Subscription。
  • Catch-up Subscription:Catch-up Subscription並不會在Server端儲存checkpoint,要由Consumer自行保存checkpoint。Consumer在連線到資料庫的時候告訴資料庫要讀取哪一個event stream,以及要從哪一個位置開始讀起,如圖5所示。

 

▲圖4:EventStoreDB的Persistent Subscription在Server端保存checkpoint。

 

▲圖5:EventStoreDB官方網站的Catch-up Subscription指定讀取位置程式範例

 

***


實作Persistent Consumer

講了這麼多,接下來要寫Persistent Consumer將checkpoint存在Message DB。因為是Event Sourcing的系統,所以在資料庫端checkpoint也會存在某個代表Consumer的event stream裡面。請參考圖6,stream_name欄位的值是$$Checkpoint-ezkanban-11,其中「$$Checkpoint-」這個前置字串代表它是一個系統產生用來儲存checkpoint的stream,ezKanban-11則是Consumer的名字。type欄位的值是$System$Checkpointed,代表它是一個產生或更新checkpoint的事件。最後,data欄位儲存 {“position”: 5},代表checkpoint的數值,目前讀到event stream第5個位置。

理論上,因為是Event Sourcing系統,所以每次更新checkpoint的數值應該是寫入一筆新的事件,然後這個checkpoint stream的最後一筆資料就是目前最新的讀取位置。在這裡Teddy採用傳統CRUD的做法,只儲存最新的一筆checkpoint資料。也就是說,更新checkpoint不會寫入一筆新的事件,而是直接更新原有事件的data欄位。

 

       

▲圖6:Checkpoint儲存在代表Consumer的event stream裡面

 

 

圖7是產生PresistentConsumer的程式碼,第50行判斷這個Consumer是否是第一次建立,如果是在第51行呼叫_writeMessage方法產生一個新的event stream並寫入一筆checkpoint=0的資料。

 

▲圖7:Checkpoint儲存在代表Consumer的event stream裡面

 

為了從event stream讀取事件,PresistentConsumer必須定期向Event Store查詢,如圖8所示。第35行到41行取得checkpoint的值,第43行執行一個while(true)迴圈,在45行從指定的checkpoint位置讀取$all stream的事件(這個Consumer的用途是用來監聽所有系統事件)。讀到資料之後就可以處理它們,處理完畢之後第51行呼叫ack方法寫入新的checkpoint位置到資料庫中。這裡有一個實作細節要注意:正常情況下第44行程式在讀取事件的時候會多設定「一次最多讀幾筆資料」的batch size參數,以免萬一event stream裡面有非常巨量的資料,造成Consumer卡住甚至當機(可能用光記憶體)的問題。

若監聽的event stream裡面沒有新的事件,則會直接sleep一段時間(polling interval)再重查一次。

 

▲圖8:PresistentConsumer的run方法

 

圖9為ack程式碼,首先判斷event stream是否存在(第76~79),以及checkpoint的位置是否大於event stream裡面最後一個事件的位置(第81~84),最後檢查stream name不是系統內建的stream。如果都沒問題,就寫入新的checkpoint(第90行)。



▲圖9:ack方法

 

***

 

下集預告

介紹完在如何在伺服器端紀錄checkpoint以達到事件至少傳遞一次,下一集介紹如何實做Idempotent。

***

友藏內心獨白:連載已經接近尾聲了。



2022年7月13日 星期三

事件溯源(16):分散式系統的事件語意與Idempotent

July 6 23:14~24:00;July 7 00:00~01:39

▲快寫到體力不支了 XD

 

前言

鄉民們之所以要採用Event Sourcing與CQRS,很多情況都是為了開發微服務。微服務架構屬於分散式系統,相較於集中式系統,分散式系統具有異質性、容易擴展、比較強健(robust)、容錯性高且系統模組之間的耦合性較低等特性。但相對地,分散式系統的開發也比較複雜且一不小心就容易「出錯」。

在DDD中,Aggregate之間的狀態透過領域事件達到最終一致性(eventual consistency)。相似地,在微服務架構下,各個微服務之間的狀態同步也是透過事件或是訊息達到最終一致性。但在分散式系統中,事件傳遞本身也可能發生遺失,導致接收者收不到事件,也就做不到最終一致性。

本集介紹在分散式系統中為了正確做到最終一致性,事件傳遞與事件處理器(event handler)須具備那些特性。

 

***

保證事件的順序

首先,事件本身的順序不能亂掉,否則事件接收者的狀態一定會出錯。當事件被寫入Event Store的當下,Event Store本身必須保證在同一個event stream裡面事件必須依據發生(寫入時間點)的先後順序排序,這點基本上沒有問題。看到這裡鄉民們可能會想:「事件在Event Store中既然已經排序過,為什麼還會發生順序亂掉的情況?」

請參考圖1,事件在Event Store中一開始順序是正確的,假設這個事件是ezKanban裡面User Management Bounded Context所產生的事件,像是UserCreated、UserRenamed、UserEmailChanged等。在ezKanban中,Kanban Board Bounded Context(ezKanban的Core Domain)需要在畫面上顯示使用者名稱,所以它會聽UserCreated、UserRenamed等領域事件,然後在自己本地端建立一份User資料的快取。因此,User Management Bounded Context會把內部的UserCreated這些領域事件往外傳,寫到Pulsar的Topic A裡面。

假設在Kanban Board Bounded Context裡面,為了「求快」啟動兩隻consumer(event handler)程式同時間去讀取Topic A的資料。這種consumer在訊息導向架構中稱為Competing Consumer(競爭消費者),它們會搶著處理Topic裡面的資料。圖1中的Consumer A和Consumer B從Topic A處理完事件之後會寫另外一筆事件到Topic B。現在Consumer A拿到了1和3這兩個訊息,Consumer B拿到了2和4這兩個訊息,因為它們執行的速度不同,所以處理完之後最後Topic B裡面事件的順序變成2, 1, 4, 3,和原始事件產生的順序不同。



 ▲圖1:Event out of order示意圖

 

如果今天Topic A存放的是轉檔需求,而Topic B存放的是轉檔完成的結果,在這種情況下通常來說事件順序並不重要。但是如果是要透過Event Broker傳遞事件然後希望接收者達到最終一致性,就要注意在事件傳遞的過程中,不要不小心造成事件順序亂掉的情況。

***

至少一次(At least once)

除了事件順序不能錯,另一個關於事件傳遞的要求,就是事件要保證至少會被接收者看到一次(at least once)

為什麼要這麼麻煩?怎麼不規定恰好一次(exactly once)就好?因為做不到。請參考圖2,Consumer把Event 1從Topic A讀出之後,它事情還沒做完就當掉。Consumer重啟之後,因為Event 1已經被讀走,Topic A裡面最新的事件變成Event 2。但是Event 1還沒被Consumer處理,也就是說Event 1從此就從地球上消失。

 

▲圖2:Consumer讀取事件之後,事情還沒做完就當掉

 

***

要解決事件消失的問題也很簡單,就是Consumer讀出事件之後,Topic不會立刻把事件刪除,一直到Consumer處理完畢並通知(ack)Topic,此時Topic才會把事件刪除,如圖3所示。在圖3中,情況1,Consumer讀出事件1後立刻當掉。但因為Consumer沒有ack,所以事件1還存在Topic A。情況2當Consumer重啟之後,又可以讀到一次事件1(這就是at least once,至少一次,至多不限)。情況3當Consumer執行完畢並且ack,事件1才會從Topic A移除。


▲圖3:Consumer通知Topic之後被讀出的事件才會從Topic刪除

 

***

 

Idempotent

現在新的問題來了,ack的作法雖然可以確保事件至少被Consumer收到一次,但當Consumer重複收到同一個事件怎麼辦?例如收到重複的扣款請求?那就真的變成「詐騙集團」常用的台詞:「系統設定錯誤造成重複扣款XD」。

請參考圖4,在情況4中Consumer收到事件1並且已經把事情做完(更新系統狀態),就在它要ack之前,它當掉了,所以沒有ack成功。情況5當它下次重啟,事件1又被處理一次,這樣顯然不OK。因此Consumer(Event Handler)需要具備Idempotent。


▲圖4:事件1被Consumer處理2次

 

Idempotent意指相同操作就算重複執行也不會影響系統狀態。例如,把任何數字乘以1,最後結果還是不變,因此「乘以1」這個操作就是idempotent。在軟體開發中,傳統的CRUD操作,基本上RUD都是idempotent。R不用說,讀取資料N次也不會改變系統狀態,用固定值更新與刪除同一筆資料N次也不會改變系統狀態。但是C(新增)就不是idempotent。

Consumer需要具備idempotent的意思,就是說Consumer可以很神奇地讓兩次相同新增的效果變成一次。怎麼做到?原則上就是讓Consumer把它讀過的事件編號(event id)記錄下來。每次從Topic讀取事件之後,先到自己本地端資料庫查詢這筆事件以前有沒有看過?如果有就直接ack繼續處理下一筆,如果沒有就正常處理,如圖5所示。

圖5中有一個細節要注意,Consumer儲存處理過的領域事件與因為處理該事件所造成的系統狀態更新,這兩的操作必須要放在同一個交易中,否則又可能會發生Consumer改變狀態但來不及紀錄事件,或是先記錄事件但是來不及保存狀態的錯誤狀況

 

▲圖5:儲存處理過的事件以達到idempotent


 

***


好多細節

看到這理請鄉民們回頭看<Consumer事件溯源(10):實作Projector>,Projector是一種Consumer,因此在實作Projector的時候就必須考慮到本集所提到event ordering、at least once以及Idempotent的問題。

▲圖6:Read Model的Projector需要具備Idempotent

 

***

 

下集預告

如果鄉民使用EventStoreDB,它本身既是Event Store也是一個簡易的Event Broker,因此支援event ordering與at least once。Kafka與Pulsar更是強大的Event Broker,當然也有支援。但是,如果鄉民是自己用Rational Database或是NoSQL「手刻Event Store」,那麼就需要自己確保event ordering與at least once。下集談在手刻Event Store的情況下,要如何做到event ordering與at least once。

***

友藏內心獨白:身為Maker一定要自幹Event Store的啦XD。

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設計模式終於派上用場了。