
系统描述设计一个分布式的cache可以被多个backend service用来缓存和访问数据。系统需求功能需求1. 系统支持- GET(key)- SET (key, value, TTL)- DELETE(key)2. 支持多个backend services共享cache3. cache可以存储任意的key/value4. 支持TTL5. 当cache miss时client可以从underlying database/service获取数据并重新写入cache6. 支持eviction例如LRU。7. 支持热点访问数据。非功能需求1. 10 billion requests/day平均116k requests/sec, peak 1M requests/sec2. 每个value平均10kB最大value1MB3. cache cluter初始容量100TB未来扩展到1PB4. GET P99 lateny 5ms5 availability99.99%6. 接受eventual consistency7. cache failure不能导致整个业务系统不可用API定义// GET data GET /data/{key} response: { key: ..., value: ... } // SET data POST /data { key: ..., value: ..., ttlseconds: 3600 } response: status 200 OK // DELETE DELETE /data/{key} response:status 200 OK根据题目的描述value的平均值为10KB最大为1MB这个data load可以直接放进request的request body本身。架构描述1. 首先查看我们的scale平均116k request/secpeak为1M requests/sec。两种traffic都很难使用单一的instance处理。假设单一instance可以支持10k request/sec则我们需要build近12个instance的cluster。2. 因为cache是为backend service使用的因此我们能不需要api gateway。我们设置client library可以用来访问cache router。cache router维护多个nodes每个node带有memory用来进行cache存储。cache router使用consistent hash算法维护key和node之间的mapping关系。3. 因为我们需要在cache miss的时候保证数据可以从DB中读出我们可以在client library提供这样的逻辑当用户尝试GET data时我们首先访问cache如果miss由cache router访问DB。假如读取到data则保存到cache中。在delete data时则同时从DB和cache中删除。update时可以先update DB然后从cache中删除stale的value。4. cache的详细设计因为我们要支持big scale datacache由多个node组成。一个key写入到哪个node计算hash value然后通过consistent hash维护key和node的对应。这样在某些node出现不可用问题时有利于减少转移数据的影响。另外我们为多个node组成的cluster创建和维护connection pool来减少访问shard所需要的latency。5. 为了进一步提高shard的availability我们可以为shard设置replica。这样既可以提高shard可以支持的scale也进一步提高了shard的availability。update cache时首先update prmary node然后primary node将数据进一步sync up到它的replica中。在GET时primary node和replica都可以进行访问。DELETE时由primary node负责notify replica删除数据然后primary node删除返回。6. 为了维护cache的evicatioin各个node在内部的memory构建LRU的chain。GET, DELETE,和UPDAE的复杂性都是O(1)。7. 为了维护TTL我们可以使用lazy evaluation的方法。也就是直到GET的时候我们再判断一个key的TTL是否已经expire。如果expire了则不返回该record而是删除该record。但是这样我们的cache中会有大量不需要的record则我们也可以使用一张额外的表以nodeId作为primary keyTTL作为sort key再设置一个cleaner定期查询每个node最早expire的TTL删除数据。也可以在node内部的memory使用heap维护这样的信息并启动一个thread执行类似的操作。7. 在cache failed的情况下我们的service会fallback到访问DB所以数据访问仍然可用。8. 当出现hot key时一个node会成为访问的瓶颈。一个办法是添加client端的cache和CDN。另一个一个办法是增加replica分散一个node instance的压力。系统讨论1. DB和cache之间的关系DB是data的source of truth。因此在我们同时需要update DB和cache时我们要保证DB被update成功而cache里的数据可以直接删除这样下一次访问时DB中的记录会自动sync up到cache里。或者我们为cache里的记录设置default TTL。这样保证cache里stale的数据在某段时间后会被自动eviction。2. 在cache失效时如何安全的访问DB在cache失效时直接访问DB很容易造成DB过载。所以我们应该设置rate limiterconnection pool limitbackpressurelocal cache来应对这种极端情况。3. 进一步优化cache的访问在有多个GET访问同一个key时我们可以在一个interval内将多个GET等待一个node GET操作。这样我们将reduce大量memory的重复GET操作。