FAN YOURSS

系统设计[实战篇 — 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的操作有几种选择。

  1. 使用db的auto increment
  2. 读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