List the key functional requirements for the system (Ask the AI for hints if stuck)...
List the key non-functional requirements (performance, scalability, reliability, etc.)...
Estimate the scale of the system. Consider daily active users, read/write ratio, storage requirements, bandwidth, and any relevant QPS calculations...
Assuming we have 100 million DAU, and on average, each user sends 5 read requests (view balances, view accounts, view history etc.), and sends 1 write request (transfer balance, link accounts etc.) per day.
Assuming that peak QPS is twice of average QPS.
Peak read QPS = 100 million * 5 / 100k secons per day * 2 = 10k
Peak write QPS = 2k
We need to store:
For user and account metadata, for 100 million users, assuming each metadata needs 2kb.
Storage needed for each use case is:
100 million * 2kb = 200GB
For payment history, assuming 100 million payments were made per day, and 1kb is needed for each event, the storage needed is:
100 million * 1kb = 100GB per day.
For events older than 1 year, we will put them in cold storage.
Define the APIs expected from the system. This is your chance to analyze and define the read and write paths so that you can come up with the high-level design...
CURD APIs on user accounts:
POST v1/register_account
POST v1/update_account
GET v1/account_info
DELETE v1/account
Link user account with external financial situations:
POST v1/link_finanical_account
POST v1/unlink_financial_account
Deposit and withdrawal:
POST v1/deposit
POSt v1/withdraw
Transfer money to peers and accept transfer from peers:
POST v1/transfer
POST v1/accept_transfer
Make payments:
POST v1/pay
View payment history:
GET v1/payment_history
Describe the overall system architecture. Identify the main components needed to solve the problem end-to-end. Use the diagramming tool to create a block diagram.
For this design, all requests first go through a load balancer and an API gateway. Load balancer distributes traffic evenly to different API gateway instances. The API gateway handles authentication, authorization and rate limiting.
The traffic between the components in the system are encrypted via TLS.
For viewing user account metadata, the request is routed to account service, which queries account database. If the user wants to see the real time balance, the account service will query the ledger database, and returns the updated balance of the user.
For linking with external financial institutions, we query the deposit/withdraw service, and links user's digital wallets to external banks, going through an external verification process while doing so. For payment/transfer/deposit etc. we update the status in ledger database. For these operations, we require strong consistency, a sample flow looks like:
Request: Jack pays merchant $100
↓
Payment Service
↓
DB transaction #1
────────────────────────
Create transaction = INITIATED
Atomic conditional update:
available_balance -= 100
pending_balance += 100
only if available >= 100
Write reservation/ledger state
Write idempotency record
COMMIT
────────────────────────
↓
release locks
↓
Call external payment provider
Then after calling the external provider:
Success:
DB transaction #2:
pending -= $100
transaction → COMPLETED
write final ledger entries
write outbox event
COMMIT
Definite failure:
DB transaction #2:
pending -= 100available+=100
available += 100available+=100
transaction → FAILED
write outbox event
COMMIT
For the ledger database, it has an outbox that process important updates, and emits to downstreams CDC. CDC triggers events on Kafka. The downstreams consumers of Kafka are:
How do we handle race conditions?
Suppose Jack has a balance of 100 to Alice and 100 balance and incorrectly succeed.
For balance deductions, we use an atomic conditional update inside an ACID database transaction:
BEGIN;
UPDATE wallets
SET available_balance = available_balance - 100
WHERE wallet_id = 'jack'
AND available_balance >= 100;
We then check the number of affected rows.
1 row updated → sufficient balance, continue
0 rows updated → insufficient balance, rollback
If the deduction succeeds, the same transaction updates the recipient's balance and writes the corresponding transaction and ledger entries:
UPDATE wallets
SET available_balance = available_balance + 100
WHERE wallet_id = 'alice';
INSERT INTO transactions (...);
INSERT INTO ledger_entries (... Jack -100 ...);
INSERT INTO ledger_entries (... Alice +100 ...);
COMMIT;
If Jack sends two 100 → $0. When the other request evaluates the condition available_balance >= 100, it fails and updates zero rows.
This prevents double spending without requiring application-level optimistic locking.
For more complex workflows where we must read several values, perform application logic, and then write them back, we could instead use optimistic locking with a version number, or pessimistic locking such as SELECT ... FOR UPDATE when contention is high.
How do we implement the offline payment flow?
True offline payment cannot provide the same real-time strong consistency as an online transaction. We either reject offline payments entirely or accept bounded double-spend risk using pre-authorized, cryptographically protected offline spending credentials.
Define the data model. Identify the main entities, their attributes, and relationships. Consider the choice of database type (SQL vs NoSQL) and justify your decision based on access patterns...
In this design, there are 3 types of data we need to store:
For all of the data, we will use distributed relational database, the benefits are:
The tradeoffs are:
These are acceptable tradeoff in our use case.
The data schemas are:
users
-----
user_id: UUID
user_name: String
password_hash: String
email: String
phone_number: String
metadata: JSON
created_at: Timestamp
wallets
-------
wallet_id: UUID
user_id: UUID
currency: String
available_balance: DECIMAL
pending_balance: DECIMAL
version: BIGINT
updated_at: Timestamp
transactions
------------
transaction_id: UUID,
idempotency_key: String,
transaction_type: ENUM
status: ENUM
created_at: Timestamp
completed_at: Timestamp?
metadata: JSON
ledger_entries
--------------
ledger_entry_id: UUID
transaction_id: UUID
wallet_id: UUID
amount: DECIMAL
currency: String
entry_type: ENUM // DEBIT / CREDIT
created_at: Timestamp
A single atomic transaction to transfer $100 between alice and bob looks like:
BEGIN;
INSERT INTO transactions (...);
INSERT INTO ledger_entries (... -100 ...);
INSERT INTO ledger_entries (... +100 ...);
UPDATE wallets
SET available_balance = available_balance - 100
WHERE wallet_id = 'bob';
UPDATE wallets
SET available_balance = available_balance + 100
WHERE wallet_id = 'alice';
COMMIT;
As data size and traffic increase, we need replication and partitioning to maintain availability, durability, and scalability.
Replication: Each database shard/range is replicated across multiple nodes using a consensus protocol such as Raft. Writes are sent to the leader and replicated to a majority of replicas before being acknowledged as committed. Strongly consistent reads can be served through the leader/leaseholder. If the leader fails, the remaining replicas can elect a new leader, as long as a majority of replicas are still available. Because every committed write has already been persisted by a majority, committed financial transactions are not lost during leader failover.
Sharding/partitioning: Wallet data is partitioned primarily by wallet_id, so common operations such as retrieving or updating one wallet's balance are normally local to a single shard/range. P2P transfers may involve wallets on different shards. We therefore use a distributed relational database such as CockroachDB that supports ACID transactions across shards, accepting the additional latency of distributed transactions for these cases.
Deep dive into 2-3 key components. Explain how they work, how they scale, discuss tradeoffs, capacity, and any relevant algorithms or data structures.
In this design, there are 3 types of data we need to store:
For all of the data, we will use distributed relational database, the benefits are:
The tradeoffs are:
These are acceptable tradeoff in our use case.
The data schemas are:
users
-----
user_id: UUID
user_name: String
password_hash: String
email: String
phone_number: String
metadata: JSON
created_at: Timestamp
wallets
-------
wallet_id: UUID
user_id: UUID
currency: String
available_balance: DECIMAL
pending_balance: DECIMAL
version: BIGINT
updated_at: Timestamp
transactions
------------
transaction_id: UUID,
idempotency_key: String,
transaction_type: ENUM
status: ENUM
created_at: Timestamp
completed_at: Timestamp?
metadata: JSON
ledger_entries
--------------
ledger_entry_id: UUID
transaction_id: UUID
wallet_id: UUID
amount: DECIMAL
currency: String
entry_type: ENUM // DEBIT / CREDIT
created_at: Timestamp
A single atomic transaction to transfer $100 between alice and bob looks like:
BEGIN;
INSERT INTO transactions (...);
INSERT INTO ledger_entries (... -100 ...);
INSERT INTO ledger_entries (... +100 ...);
UPDATE wallets
SET available_balance = available_balance - 100
WHERE wallet_id = 'bob';
UPDATE wallets
SET available_balance = available_balance + 100
WHERE wallet_id = 'alice';
COMMIT;
As data size and traffic increase, we need replication and partitioning to maintain availability, durability, and scalability.
Replication: Each database shard/range is replicated across multiple nodes using a consensus protocol such as Raft. Writes are sent to the leader and replicated to a majority of replicas before being acknowledged as committed. Strongly consistent reads can be served through the leader/leaseholder. If the leader fails, the remaining replicas can elect a new leader, as long as a majority of replicas are still available. Because every committed write has already been persisted by a majority, committed financial transactions are not lost during leader failover.
Sharding/partitioning: Wallet data is partitioned primarily by wallet_id, so common operations such as retrieving or updating one wallet's balance are normally local to a single shard/range. P2P transfers may involve wallets on different shards. We therefore use a distributed relational database such as CockroachDB that supports ACID transactions across shards, accepting the additional latency of distributed transactions for these cases.
For this design, all requests first go through a load balancer and an API gateway. Load balancer distributes traffic evenly to different API gateway instances. The API gateway handles authentication, authorization and rate limiting.
The traffic between the components in the system are encrypted via TLS.
For viewing user account metadata, the request is routed to account service, which queries account database. If the user wants to see the real time balance, the account service will query the ledger database, and returns the updated balance of the user.
For linking with external financial institutions, we query the deposit/withdraw service, and links user's digital wallets to external banks, going through an external verification process while doing so. For payment/transfer/deposit etc. we update the status in ledger database. For these operations, we require strong consistency, a sample flow looks like:
Request: Jack pays merchant $100
↓
Payment Service
↓
DB transaction #1
────────────────────────
Create transaction = INITIATED
Atomic conditional update:
available_balance -= 100
pending_balance += 100
only if available >= 100
Write reservation/ledger state
Write idempotency record
COMMIT
────────────────────────
↓
release locks
↓
Call external payment provider
Then after calling the external provider:
Success:
DB transaction #2:
pending -= $100
transaction → COMPLETED
write final ledger entries
write outbox event
COMMIT
Definite failure:
DB transaction #2:
pending -= 100available+=100
available += 100available+=100
transaction → FAILED
write outbox event
COMMIT
For the ledger database, it has an outbox that process important updates, and emits to downstreams CDC. CDC triggers events on Kafka. The downstreams consumers of Kafka are:
How do we handle race conditions?
Suppose Jack has a balance of 100 to Alice and 100 balance and incorrectly succeed.
For balance deductions, we use an atomic conditional update inside an ACID database transaction:
BEGIN;
UPDATE wallets
SET available_balance = available_balance - 100
WHERE wallet_id = 'jack'
AND available_balance >= 100;
We then check the number of affected rows.
1 row updated → sufficient balance, continue
0 rows updated → insufficient balance, rollback
If the deduction succeeds, the same transaction updates the recipient's balance and writes the corresponding transaction and ledger entries:
UPDATE wallets
SET available_balance = available_balance + 100
WHERE wallet_id = 'alice';
INSERT INTO transactions (...);
INSERT INTO ledger_entries (... Jack -100 ...);
INSERT INTO ledger_entries (... Alice +100 ...);
COMMIT;
If Jack sends two 100 → $0. When the other request evaluates the condition available_balance >= 100, it fails and updates zero rows.
This prevents double spending without requiring application-level optimistic locking.
For more complex workflows where we must read several values, perform application logic, and then write them back, we could instead use optimistic locking with a version number, or pessimistic locking such as SELECT ... FOR UPDATE when contention is high.
How do we implement the offline payment flow?
True offline payment cannot provide the same real-time strong consistency as an online transaction. We either reject offline payments entirely or accept bounded double-spend risk using pre-authorized, cryptographically protected offline spending credentials.
Also before we fully commit the payment and return success to users, we can add a fraud screening service that uses ML models to flag suspicious activities, and alert using notification service in case of positive findings.
For the payment gateway to query external payment providers, we want to put a circuit breaker. So external provider failures don't cause cascading errors on our side.