系统设计[实战篇 — Snowflake]
今年的第一个系统面试告负,拖了一个多星期终于静下心来准备写下面试的经历,一方面失败总是难以接受,另一方面,这次的面试对自己仿佛也有了更新的认识。长期在当前职位的react而不是think,也仿佛导致了必然的结果。
Snowflake的面试官都是属于人狠话不多类型,上来直接丢题,对你的背景毫不感兴趣。题目是设计Key-Value Store with global versioning and time travel。听到Global Key-Value Store,我从Design Data-Intensive Application书里学来的东西直接short circuit了我的大脑,按着pattern就一点一点开始设计Key Value Store,直接完全忽略了global versioning 和 time travel的requirement,导致最后压根都涉及到time travel,完全纠结在make the basic work。
Redesign
应该从data出手,想想in-memory和single db如何handle,以及对performance的需求,然后延展到distributed system.
比如这里的核心需求是global versioning,那么所有的历史数据都必须得存下来。同时有time travel的需求,那么timestamp也必须得被存下来。
那么最简单的一个table schema就可以变成:
-- create_snapshot_table.sql
CREATE TABLE IF NOT EXISTS snapshots (
snapshot_id SERIAL PRIMARY KEY,
key VARCHAR NOT NULL,
value VARCHAR NOT NULL,
created_at TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP
);
CREATE INDEX IF NOT EXISTS idx_snapshots_key ON snapshots (key);
CREATE INDEX IF NOT EXISTS idx_snapshots_created_at ON snapshots (created_at);
每个put request会append一个row,每一个row number都是snapshot id。
这个本质上跟SSTable一样,只不过不涉及到底层的文件存储,这些都被postgres 的api abstract掉了。
这样put和get包括ss和timetravel的handling就是简单的SQL.
Read Performance — page miss
一旦page miss,比如以下query,read=2,表示有两次buffer以外的read。
因为本机的postgres是跑在ssd上的,所以一次的read大概是1ms,正好呼应jeff dean的 latency number: https://gist.github.com/jboner/2841832
key-value-store=# EXPLAIN (ANALYZE, BUFFERS) SELECT value FROM snapshots WHERE key = '2U' AND created_at <= '2023-11-06' ORDER BY created_at DESC LIMIT 1;
QUERY PLAN
---------------------------------------------------------------------------------------------------------------------------------------
Limit (cost=23.98..23.98 rows=1 width=48) (actual time=2.010..2.011 rows=1 loops=1)
Buffers: shared hit=25 read=2
-> Sort (cost=23.98..23.99 rows=5 width=48) (actual time=2.008..2.009 rows=1 loops=1)
Sort Key: created_at DESC
Sort Method: top-N heapsort Memory: 25kB
Buffers: shared hit=25 read=2
-> Bitmap Heap Scan on snapshots (cost=4.46..23.95 rows=5 width=48) (actual time=1.968..1.995 rows=24 loops=1)
Recheck Cond: ((key)::text = '2U'::text)
Filter: (created_at <= '2023-11-06 00:00:00'::timestamp without time zone)
Heap Blocks: exact=24
Buffers: shared hit=25 read=2
-> Bitmap Index Scan on idx_snapshots_key (cost=0.00..4.46 rows=5 width=0) (actual time=1.952..1.953 rows=24 loops=1)
Index Cond: ((key)::text = '2U'::text)
Buffers: shared hit=1 read=2
Planning:
Buffers: shared hit=5
Planning Time: 0.172 ms
Execution Time: 2.054 ms
(18 rows)
对比一个没有disk read的query,execution只用了0.121ms.
key-value-store=# EXPLAIN (ANALYZE, BUFFERS) SELECT value FROM snapshots WHERE key = 'TX' AND created_at <= '2023-11-06' ORDER BY created_at DESC LIMIT 1;
QUERY PLAN
----------------------------------------------------------------------------------------------------------------------------------------
Limit (cost=78.79..78.79 rows=1 width=18) (actual time=0.097..0.098 rows=1 loops=1)
Buffers: shared hit=26
-> Sort (cost=78.79..78.84 rows=21 width=18) (actual time=0.096..0.097 rows=1 loops=1)
Sort Key: created_at DESC
Sort Method: top-N heapsort Memory: 25kB
Buffers: shared hit=26
-> Bitmap Heap Scan on snapshots (cost=4.46..78.68 rows=21 width=18) (actual time=0.038..0.085 rows=25 loops=1)
Recheck Cond: ((key)::text = 'TX'::text)
Filter: (created_at <= '2023-11-06 00:00:00'::timestamp without time zone)
Heap Blocks: exact=24
Buffers: shared hit=26
-> Bitmap Index Scan on idx_snapshots_key (cost=0.00..4.45 rows=21 width=0) (actual time=0.029..0.029 rows=25 loops=1)
Index Cond: ((key)::text = 'TX'::text)
Buffers: shared hit=2
Planning:
Buffers: shared hit=4
Planning Time: 0.163 ms
Execution Time: 0.121 ms
Read Performance — Index Scan
有意思的是,Read的performance跟这个entry多久没有update有关。一个刚update过的entry read只需要0.035ms,而一个很久没有update过的entry read需要42ms.
仔细研究了一下,居然发现问题出现在Query Planner,他没有用Binary Search,而对index 进行sequential backwards scan。
key-value-store=# EXPLAIN (ANALYZE, BUFFERS) SELECT value FROM snapshots WHERE key = 'kak' AND created_at <= '2023-11-06' ORDER BY snapshot_id DESC LIMIT 1;
QUERY PLAN
---------------------------------------------------------------------------------------------------------------------------------------------------
Limit (cost=0.42..5.69 rows=1 width=44) (actual time=42.000..42.001 rows=1 loops=1)
Buffers: shared hit=4552
-> Index Scan Backward using snapshots_pkey on snapshots (cost=0.42..11405.00 rows=2165 width=44) (actual time=41.998..41.998 rows=1 loops=1)
Filter: ((created_at <= '2023-11-06 00:00:00'::timestamp without time zone) AND ((key)::text = 'kak'::text))
Rows Removed by Filter: 241861
Buffers: shared hit=4552
Planning:
Buffers: shared hit=5
Planning Time: 0.201 ms
Execution Time: 42.072 ms
(10 rows)
问题出现在这个LIMIT 1 和 SELECT value上,再看一眼这个query:
SELECT value
FROM snapshots
WHERE key = 'kak' AND created_at <= '2023-11-06'
ORDER BY snapshot_id DESC
LIMIT 1;
value是没有index的,导致query planner认为不需要使用index。
LIMIT 1导致query planner认为从后往前scan,只要找到第一个符合filter条件return就行。
这两个条件combine,导致query planner错认为scan是最快的。
(目前尝试了各种改写query的办法都以失败告终)
最后直接去掉了LIMIT 1,让query返回所有结果,在Application layer只取第一个result。
Before:
fanyou@mbp15 key-value-store % ab -n 1000 -c 2 http://127.0.0.1:8000/kak
Percentage of the requests served within a certain time (ms)
50% 111
66% 111
75% 113
80% 113
90% 114
95% 123
98% 130
99% 141
100% 142 (longest request)
After:
fanyou@mbp15 key-value-store % ab -n 1000 -c 2 http://127.0.0.1:8000/kak
Percentage of the requests served within a certain time (ms)
50% 17
66% 17
75% 17
80% 17
90% 20
95% 27
98% 29
99% 31
100% 37 (longest request)
Read Performance — Scaling
最后整个db对每个entry的query P50大概在1.7ms 左右,server response time P50大概在9ms左右。
这些都是可以通过砸更多机器保证的。
进一步的提升那么只能通过cache,可以达到sub-ms 的read速度 (100ns),或者干脆在client端cache,从而完全不需要network request。
Write Performance — Snapshot id
对Snapshot id的操作有几种选择。
- 使用db的auto increment
- 读current ss然后写入ss+1,利用unique index constraint on primary key去guarantee consistency
2 显然是没法scale的,在很多concurrent write发生的时候,大家都会争夺current ss id,然后只有一个可以成功, 其他所有都要retry。
1会快很多。
Write Performance — Scaling
最后Write的performance大概在 P50~2ms.
有意思的是对于已经insert过的key,再次insert P50会只有 ~0.2ms,快了10倍,目前猜测是因为index无需insert。
Write 的scaling是这里的重头,假设recurring key很多,那么理论RPS大概是~5000
如果new key很多,那么理论RPS大概是~500
也就是最多支持5000个sensor不停的写入DB。
Scaling最难的地方大概在于maintain 这个global snapshot id,一旦开始scale write,基本就要每个DB node都能写,那么就不能避免的会有conflict resolution,而不能简单的靠increment id了。
Idea1: [No Table Schema Change] Postgres Multi Master Setup
???
Idea2: [No Table Schema Change] Switch to Cassandra, W + R ? N (tunable consistency)
??? 不见得会有提升,因为大家都还是需要rely on这个single id,either lose data或者lose performance。
Idea 3: [Change Table Schema] Separate Snapshot Tracking from Data
snapshot id -> uuid
每个write还是生成一个snapshot id (uuidv4),根据snapshot id找到timestamp,然后用timestamp往前搜索。
这样就enable sharding,同时可以用Cassandra之流。
Read会慢一些,因为timestamp的sort会span across multiple database shard。
但是write也可以horizontal scale了。
Idea 4: Batch Snapshot Id Update
每个Write都先进入一个Event Queue,然后Worker进行batch snapshot update,这样write没法立刻得到Snapshot id。不符合需求。
Failure Handling
Request -> Coordinator (API Service) -> Persistence (Postgres)
- When write is received, but Coordinator crashed, write is lost
- other failures? if postgres is down, cannot read/write, add cache so we can still support read
Reference
https://github.com/Noeyfan/system-design/tree/master/key-value-store