Wrap Up
In this chapter, we designed a wallet service that is capable of processing over 1 million payment commands per second. After a back-of-the-envelope estimation, we concluded that a few thousand nodes are required to support such a load.
In the first design, a solution using in-memory key-value stores like Redis is proposed. The problem with this design is that data isn't durable.
In the second design, the in-memory cache is replaced by transactional databases. To support multiple nodes, different transactional protocols such as 2PC, TC/C, and Saga are proposed. The main issue with transaction-based solutions is that we cannot conduct a data audit easily.
Next, event sourcing is introduced. We first implemented event sourcing using an external database and queue, but it’s not performant. We improved performance by storing command, event, and state in a local node.
A single node means a single point of failure. To increase the system reliability, we use the Raft consensus algorithm to replicate the event list onto multiple nodes.
The last enhancement we made was to adopt the CQRS feature of event sourcing. We added a reverse proxy to change the asynchronous event sourcing framework to a synchronous one for external users. The TC/C or Saga protocol is used to coordinate Command executions across multiple node groups.
Congratulations on getting this far! Now give yourself a pat on the back. Good job!
Chapter Summary
Reference materials
- Transactional guarantees: https://docs.oracle.com/cd/E17275_01/html/programmer_reference/rep_trans.html
- TPC-E Top Price/Performance Results: http://tpc.org/tpce/results/tpce_price_perf_results5.asp?resulttype=all
- ISO 4217 CURRENCY CODES: https://en.wikipedia.org/wiki/ISO_4217
- Apache Zookeeper: https://zookeeper.apache.org/
- Martin Kleppmann (2017). Designing Data-Intensive Applications. O'Reilly Media.
- X/Open XA: https://en.wikipedia.org/wiki/X/Open_XA
- Compensating transaction: https://en.wikipedia.org/wiki/Compensating_transaction
- SAGAS, HectorGarcia-Molina: https://www.cs.cornell.edu/andru/cs711/2002fa/reading/sagas.pdf
- Evans, E. (2003). Domain-Driven Design: Tackling Complexity in the Heart of Software. Addison-Wesley Professional.
- Apache Kafka. https://kafka.apache.org/
- CQRS: https://martinfowler.com/bliki/CQRS.html
- Comparing Random and Sequential Access in Disk and Memory: https://deliveryimages.acm.org/10.1145/1570000/1563874/jacobs3.jpg
- mmap: https://man7.org/linux/man-pages/man2/mmap.2.html
- SQLite: https://www.sqlite.org/index.html
- RocksDB. https://rocksdb.org/
- Apache Hadoop: https://hadoop.apache.org/
- Raft: https://raft.github.io/
- Reverse proxy: https://en.wikipedia.org/wiki/Reverse_proxy
Finished reading?
Mark it complete to track your progress.