Here is the link for one hour video.
From the 2010 USENIX Annual Technical Conference Keynote Address, Bobby Johnson, Director of Engineering, Facebook, Inc. discusses how in just over six years Facebook has grown from an idea in a dorm room to one of the most visited sites on the Internet. This explosive growth has created enormous technical challenges. He talks about some specific technical challenges Facebook has faced and the general principles they employ when addressing problems of scale. He also discusses how they structure their engineering process and culture to stay on top of unceasing growth while still moving fast to build new products.
20:00
Memcached vs my SQL
graph data - just right
Set operations
Append
Remove
Intersect
Client performance
34:00
Control and responsibility go together
Centralized - a variety of model - things can be done faster, gate keeper, distributed out.
Hybrid model - one person knows production team
Make things a lot of faster.
38:00
Lessons learned
Federate everything
Keep failures contained
Dark launch new features
New code is slow
Give developers room to try things
Nobody's job is to say no
Questions
40:00
1. Engineer culture?
Feel responsible but no control
2. TCP, application level, socket , Linux, what change?
Server side, custom kernel?
3. Memcached, all data are in memcached?
Photo id, ...
video, in disk
Web server -> data
4. New code is slow?
PHP profile -> work load
5. Data consistency?
6. Roll out new feature? One week to find out about memory leak, one instance mentioned
Answer:
7. Long term data storage?
From January 2015, she started to practice leetcode questions; she trains herself to stay focus, develops "muscle" memory when she practices those questions one by one. 2015年初, Julia开始参与做Leetcode, 开通自己第一个博客. 刷Leet code的题目, 她看了很多的代码, 每个人那学一点, 也开通Github, 发表自己的代码, 尝试写自己的一些体会. She learns from her favorite sports – tennis, 10,000 serves practice builds up good memory for a great serve. Just keep going. Hard work beats talent when talent fails to work hard.
Thursday, August 8, 2019
The Storage Technologies Behind Facebook Messages
Here is the link.
The engineering team behind Facebook Messages spent the past year building out a robust, scalable infrastructure. We spent a few weeks setting up a test framework to evaluate clusters of MySQL, Apache Cassandra, Apache HBase, and a couple of other systems. This talk looks at why we ultimately chose HBase and a number of the infrastructure decisions the team made to handle well over 100 billion messages per month.
The engineering team behind Facebook Messages spent the past year building out a robust, scalable infrastructure. We spent a few weeks setting up a test framework to evaluate clusters of MySQL, Apache Cassandra, Apache HBase, and a couple of other systems. This talk looks at why we ultimately chose HBase and a number of the infrastructure decisions the team made to handle well over 100 billion messages per month.
Apache HBase
Here is content to read related to Apache HBase.
Features
- Linear and modular scalability.
- Strictly consistent reads and writes.
- Automatic and configurable sharding of tables
- Automatic failover support between RegionServers.
- Convenient base classes for backing Hadoop MapReduce jobs with Apache HBase tables.
- Easy to use Java API for client access.
- Block cache and Bloom Filters for real-time queries.
- Query predicate push down via server side Filters
- Thrift gateway and a REST-ful Web service that supports XML, Protobuf, and binary data encoding options
- Extensible jruby-based (JIRB) shell
- Support for exporting metrics via the Hadoop metrics subsystem to files or Ganglia; or via JMX
Large-Scale Low-Latency Storage for the Social Network - Data@Scale
Here is the link.
Harrison Fisk, engineering manager at Facebook, talks about the storage considerations and challenges in serving Facebook's graph and includes an overview of TAO, Facebook's distributed graph storage system built on top of MySQL, and Wormhole, a publish-subscribe system.
Data Focused
Messages
. Large historical dataset
. Mostly append data
. Private data
.Low amount of viewers
.Hbase + caches
What is Hbase? 5 minutes break here to look up
Titan server - 10 minutes break here to look up
Mobile WWW
cell
Titan server
HBase
HDFS
Timeline
What is consideration of Timeline
.Lots of writes
.Single user reads
.Read recent data
.Large range reads of older data
.MySQL + Aggregator + MC
Graph data
Graph data
.Low MS response times
. <1 MS average
.Near 100% correct
.Heavy read-to-write ratio
.99.99%+reads
.Moving fast
Conclusion
.Applications are increasingly data-driven
.Different systems for different needs
.Clearly identify your goals
.Build towards those goals
graph serving
Why MYSQL?
Compression
Range query
TAO
.Distributed graph storage system
.Read and write through
.Understands nodes and edges
.Complex graph logic
.Easy to define objects/associations
TAO read path
TAO write path - local region
Master region - front end cluster
15:28
Consistency model
. Read-after-write
.Write-through
.Sticky user to cluster
.Always progression forward
.Single master
.Eventual cache consistency
Wormhole
.Pub-Sub system
.Subscribe to certain types of data change events
.All writes to a certain assoc type
.All deletes of objects
.Tails MySQL binlog consuming SQL comments
.Other systems can be tailed as well.
Wormhole consumers
.Cache invalidation
.Indexing
.Graph search
.Secondary index services
.ETL to data warehouse
Conclusion
.Social graph data needs worldwide realtime access
.Distributed caching can hide latency issues
.But create consistency issues
.Tao/MySQL/Wormhole is our solution
Actionable Items
I need to google search storage for system design.
Harrison Fisk, engineering manager at Facebook, talks about the storage considerations and challenges in serving Facebook's graph and includes an overview of TAO, Facebook's distributed graph storage system built on top of MySQL, and Wormhole, a publish-subscribe system.
Data Focused
Messages
. Large historical dataset
. Mostly append data
. Private data
.Low amount of viewers
.Hbase + caches
What is Hbase? 5 minutes break here to look up
Titan server - 10 minutes break here to look up
Mobile WWW
cell
Titan server
HBase
HDFS
Timeline
What is consideration of Timeline
.Lots of writes
.Single user reads
.Read recent data
.Large range reads of older data
.MySQL + Aggregator + MC
Graph data
Graph data
.Low MS response times
. <1 MS average
.Near 100% correct
.Heavy read-to-write ratio
.99.99%+reads
.Moving fast
Conclusion
.Applications are increasingly data-driven
.Different systems for different needs
.Clearly identify your goals
.Build towards those goals
graph serving
Why MYSQL?
Compression
Range query
TAO
.Distributed graph storage system
.Read and write through
.Understands nodes and edges
.Complex graph logic
.Easy to define objects/associations
TAO read path
TAO write path - local region
Master region - front end cluster
15:28
Consistency model
. Read-after-write
.Write-through
.Sticky user to cluster
.Always progression forward
.Single master
.Eventual cache consistency
Wormhole
.Pub-Sub system
.Subscribe to certain types of data change events
.All writes to a certain assoc type
.All deletes of objects
.Tails MySQL binlog consuming SQL comments
.Other systems can be tailed as well.
Wormhole consumers
.Cache invalidation
.Indexing
.Graph search
.Secondary index services
.ETL to data warehouse
Conclusion
.Social graph data needs worldwide realtime access
.Distributed caching can hide latency issues
.But create consistency issues
.Tao/MySQL/Wormhole is our solution
Actionable Items
I need to google search storage for system design.
Summer time and weight control project
August 8, 2019
It is my health research. I have to lose weight and the way I can do that is to go out play tennis, take a walk, road trip etc.
I like to write the story how I did to lose 7 lbs from 192 lb from last three or 4 weeks.
Introduction
It is my health research. I have to lose weight and the way I can do that is to go out play tennis, take a walk, road trip etc.
Stories from 192 lb to 185 lb
I like to write the story how I did to lose 7 lbs from 192 lb from last three or 4 weeks.
Serving a Billion Personalized News Feeds
Here is the link.
Feed ranking's goal is to provide people with over a billion personalized experiences. We strive to provide the most compelling content to each person, personalized to them so that they are most likely to see the content that is most interesting to them. Similar to a newspaper, putting the right stories above the fold has always been critical to engaging customers and interesting them in the rest of the paper. In feed ranking, we face a similar challenge, but on a grander scale. Each time a person visits, we need to find the best piece of content out of all the available stories and put it at the top of feed where people are most likely to see it. To accomplish this, we do large-scale machine learning to model each person, figure out which friends, pages and topics they care about and pick the stories each particular person is interested in. In addition to the large-scale machine learning problems we work on, another primary area of research is understanding the value we are creating for people and making sure that our objective function is in alignment with what people want.
Here is the link of VP of engineering.
Machine learning in news feed
- how to get the score, rank of all news
11:20/ 39:49
Model training phase 1 - boosted tree
. Feed has >100K dense features.
.First, prune these to top ~2K
. Training limited by NUM_ROWS * NUM_FEATURES
. Start with 100K features, max rows, keep most important 10K. train 10x rows, repeat
Do this for each feed events, train many forests.
Dispersion
Who is most important people in your life?
Suggesting friends of friends
I like to figure out how hard it is for me to understand this machine learning algorithm. I think that it is a good idea for me to explore as many topics as I can to prepare onsite interview from Facebook.
Feed ranking's goal is to provide people with over a billion personalized experiences. We strive to provide the most compelling content to each person, personalized to them so that they are most likely to see the content that is most interesting to them. Similar to a newspaper, putting the right stories above the fold has always been critical to engaging customers and interesting them in the rest of the paper. In feed ranking, we face a similar challenge, but on a grander scale. Each time a person visits, we need to find the best piece of content out of all the available stories and put it at the top of feed where people are most likely to see it. To accomplish this, we do large-scale machine learning to model each person, figure out which friends, pages and topics they care about and pick the stories each particular person is interested in. In addition to the large-scale machine learning problems we work on, another primary area of research is understanding the value we are creating for people and making sure that our objective function is in alignment with what people want.
Here is the link of VP of engineering.
Machine learning in news feed
- how to get the score, rank of all news
11:20/ 39:49
Model training phase 1 - boosted tree
. Feed has >100K dense features.
.First, prune these to top ~2K
. Training limited by NUM_ROWS * NUM_FEATURES
. Start with 100K features, max rows, keep most important 10K. train 10x rows, repeat
Do this for each feed events, train many forests.
Dispersion
Who is most important people in your life?
Suggesting friends of friends
Actionable Items
I like to figure out how hard it is for me to understand this machine learning algorithm. I think that it is a good idea for me to explore as many topics as I can to prepare onsite interview from Facebook.
How Twitter Uses Redis to Scale
How Twitter Uses Redis to Scale, here is the blog written in Chinese.
Chinese blog on weixin is here.
English version, the link is here on highscalability.com.
21:00/ 1:14:50
Build a redis cluster
自2010年,Yao Yu已经效力于Twitter的缓存团队。而本文主要基于她近日发表的“Scaling Redis at Twitter”演讲,主要谈Twitter的Redis扩展,同时也不局限于Redis。从演讲中不难发现,Twitter在缓存服务打造上积累了相当丰富的经验,就如你所想,Twitter使用了大量的缓存。
Timeline服务(一个数据中心)Hybrid List使用情况:
- 分配40TB左右的内存堆栈
- 3000万QPS(query per second)
- 超过6000个实例
BTtree(一个数据中心)使用状态:
- 分配65TB的内存堆栈
- 900万QPS
- 超过4000个实例
下文将会带你详细的学习BTree和Hybrid。
几个值得关注的点:
- Redis的表现非常不错,因为它将大量服务器中未使用的资源整合成有价值的服务。
- Twitter通过两个专为其业务设计的新数据模型来定制化Redis,因此他们获得了理想的性能,但是也因此受限于旧的代码,无法迅速添加新特性。因此我(Todd)一直在想,为什么他们会使用Redis来做这样的事情。只是想基于自己数据结构建立一个Timeline服务?Redis真的适合干这样的事情?
- 在网络饱和之前,分析大量节点上的日志数据,使用本地的CPU。
- 如果想获得非常高的性能,那么必须做快慢分离,数据永远都是快的代言,而命令和控制则代表了慢。
- 当下,Twitter正在使用Mesos作为作业调度程序以迁移到一个容器环境,这个做法很新颖,因此如何实现是一大看点。当然这个途径也存在弊端,比如在复杂的运行时环境指定硬件资源的使用限制。
- 中央集群管理器来监控集群。
- 缓慢的JVM,快速的C,因此他们使用C/C++重写缓存代理。
重点关注以上技术点,下面一起来看Twitter的Redis使用之道:
为什么使用Redis?
- 使用Redis驱动Timeline,也是Twitter系统内最重要的服务。Timeline是Tweet基于id的一个索引,Home Timeline则是所有Tweet连接起来形成的一个列表。User Timeline则由用户生成的Tweet组成,它们形成了另外一个列表。
- 为什么使用Redis代替Memcache?主要因为网络带宽问题(Network Bandwidth Problem)和长通用前缀问题(Long Common Prefix Problem)。
- 网络带宽问题
- Memcache在Timeline上的表现并没有Redis好,最大问题发生在fanout(推送)上。
- Twitter的读和写往往以增量方式进行,虽然每次的更新很少,但是Timeline本身的体积很大。
- 当一个Tweet产生时,它会被写入对应的Timeline中。对于某些数据结构来说,Tweet是组成它的一小部分数据。在读的时候,小批量Tweet被加载,而在滚轮向下滚动时,则会加载另外的一小批量。
- Home Timeline可能会非常大,但是用户仍然希望在同一个集合中读取,它可能会包含3000个实体。因此,在性能的优化上,避免数据库读取将非常重要。
- 增量读写使用了一个“读-改-写”的过程,因此对Timeline这样的大型对象进行small reads开销是非常昂贵的,通常会造成网络瓶颈。
- 在每秒10万+读和写的gigalink上,如果对象的平均大小超过1K,网络将成为瓶颈。
- 长通用前缀问题(其实是两个问题)
- 在数据格式上使用了一个灵活的模式,每个对象都有不同的属性组成。对每个属性都建独立的键,这需求给每个属性单独发送请求,而不是对所有的属性都进行缓存。
- 使用不同的时间戳跟踪度量。如果独立缓存每个度量样本,那么长通用前缀会被不停的重复缓存。
- 平衡度量和灵活的数据模式,使用分层的key space更令人满意。
- 通过CPU使用情况配置专门的集群。举个例子,内存键值存储对CPU的利用很低。服务器上1%的CPU时间也许能支撑1K+的RPS,基于不同数据模式会得出不同的结果。
- Reids的表现非常不错。它能看到服务器可以做什么,而不是正在做什么。对于简单的键值存储,像Redis这样的服务在服务器上可能会存在很多的CPU动态余量。
- Twitter首次为Timeline服务配置Redis是在2010年,同时搭载Redis的还有Ads服务。
- Twitter并没有使用Redis的磁盘特性。这很大程度因为在Twitter的系统中,缓存和存储都在不同的团队完成,他们会根据自己的使用来定制。也就是,对比Redis,存储团队有更好的服务。
- 热键是一个必须解决的问题,因此Twitter建立一个分层式缓存,客户缓存会自动的缓存热键。
Hybrid List
- 为Redis添加Hybrid List以获得更可预期的内存性能。
- Timeline是Tweet ID组成的列表,因此是一堆的整型,每个ID都很小。
- Redis支持两个列表类型:ziplist和linklist。Ziplist更具备空间效率,linklist则更加灵活,在双向链表下每个建会占用两个指针,对比ID的体积来说,这个开销非常大。
- 如果从内存的使用效率上看,ziplists是唯一之选。
- Ziplist的最大阈值被设置为Timeline的大小。因此在调整Ziplist的阈值之前,永远都不会储存比它更大的Timeline。这同样意味着一个产品决策,在Timeline中储存多少个Tweet。
- 对ziplist做添加和删除操作效率是非常低的,特别在列表非常大的情况下。从ziplist中使用memmove将数据删除,这里必须保证列表仍然是连续的。向ziplist中添加则需要一个内存realloc调用,以保证新实体有足够的空间。
- 鉴于Timeline的体积,写入操作存在潜在的高延时。Timeline在体积上的变化很大,大部分的用户不会频繁发送Tweet,因此他们的User Timeline都很小。Home Timeline,特别涉及到那些名流人物则可能很大。当修改一个巨大的Timeline,而内存又被缓存占满,通常会进行这个过程:大量小体积timeline将会被驱逐,直到有足够的连续内存来储存这个巨大的ziplist。这些内存管理操作将花费很多的时间,因此写入操作可能存在非常高的潜在延时。
- 因为写入会造成很多Timeline的修改操作,基于需要在内存中扩展Timeline,这里有很大的可能会产生写延时陷阱。
- 基于写延时带来的高可变性,很难为写入建立SLA。
- Hybrid List是ziplists组成的linklist,因此存在一个基于所有ziplist最大体积的阈值(以字节为单位)。以字节为单位最大程度上是为了内存效率,它可以帮助分配和解除分配相同体积的内存。当一个列表结束后,空间就被分配给下一个Ziplist。Ziplist并不可回收,直到这个列表为空。这就意味,通过删除,你可以让每个ziplist只包含一个实体。通常情况下,Tweet并不会被全部删除。
- 在Hybrid List之前,解决方案是尽快的让一个大的Timeline过期,这会给其他Timeline节约出内存,但是如果用户再查看这个Timeline时,开销是非常大的。
BTree
- 将BTree添加到Redis是为了支持分层键上的范围查询,从而得到一个结果列表。
- 在Redis中,增加次关键字或字段一般通过Hash Map处理,排序数据以执行一个范围查询时,sorted set被使用。因为sorted set只能使用一个double类型的score来排序,所以这里不能任意指定次级键或者名称来排序。同时,鉴于Hash Map使用的线性搜索,因此它并不适用于存在太多次关键字或字段的情况。
- 通过BSD实现BTree,并将之添加到Redis中。支持键查询和范围查询,同时还具备了良好的查询性能及简单的代码。唯一的缺点是BTree并不具备一个良好的内存利用效率,因为指针将造成大量的元数据开销。
集群管理
- 为每个目的建立单独的集群,每个集群部署了1个以上的Redis实例。如果一个数据集的大小大于单Redis实例可以支撑的极限,或者单Redis实例并不能提供足够的吞吐量,key space需要被分割,数据则会横跨一组实例在多个分片上保存,路由器将会为key选择应该保存的数据分片。
- 集群管理是Redis不会被撑爆的首要原因。既然可以使用集群,那么没理由不将所有的用例转移到Redis上。
- Redis集群运营起来并不轻松。基于频繁的更新操作,人们使用Redis,但是许多Redis运营并不是幂等的。比如,网络故障可能会引起重试,从而在一定程度上破坏数据的完整性。
- Redis集群支持集中管理者操控全局。使用memcache,许多集群使用了一个基于一致性哈希的客户端解决方案。如果出现非一致性数据,只能顺其自然。为了提供更好的服务,集群需要一个功能负责检测哪个分片发生问题,并且重新进行操作来保证同步。在足够长的时间后,缓存应该被清理。在Redis中破坏数据是不容易被发现的,当有一个列表,它缺失了一个部分,这个问题很难被描述清楚。
- 在建立Redis集群上,Twitter做过了很多尝试。Twemproxy,最初并不是在Twitter内部使用,它被建立用于服务Twemcache,随后添加了对Redis的支持。同时,还针对代理类型路由建立了两个额外的解决方案,其中一个与Timeline服务有关,但不是通用的;另一个则是专为Timeline设计的通用解决方案,提供了集群管理、复制和分片修复。
- 集群中存在3个选择,服务器相互间通信以达成一个协议:集群状态;使用一个代理;或者是当客户端数量达到阈值时做一个客户端方面的集群管理。
- 没有做服务器方面的优化,因为一直以保持服务器简单、透明和快速为理念。
- 并没有通过客户端,因为改变不容易被推广。在Twitter,1个缓存集群大约为100个项目使用。如果在一个客户端中做改变必须推进到100个客户端,这花费的时间可能以年计算。快速迭代意味着客户端不能放任何代码。
- 使用一个代理模式路由途径以及分片主要基于两个原因。首先,缓存服务必须是个高性能服务。如果你想获得高性能,就必须分离快快慢路径,快对应着数据路径,而慢则对应了命令行和控制路径。其次,如果将集群管理混合到服务器中,将增加Redis编码的复杂度,从而造成一个有状态的服务,如果你想给管理代码打补丁或者升级,有状态的服务你必须要重启,从而可能会导致丢失一部分数据,滚动重启集群是非常痛苦的。
- 使用代理途径的另一个原因是在客户端和服务器之间插入一个新的网络跃点。关于在添加额外的网络跃点上,Profiling 揭露了业内普遍存在的一个谣言。最起码在Twitter的系统中,Redis服务器产生的延时不会超过0.5毫秒。在Twiiter,大部分的后端系统都基于Java,并使用Finagle做相互间的交互。在通过Finagle后,延时增加到10毫秒。因此附加的网络跃点并不是问题,问题的本身在于JVM。除了JVM之外,你可以做任何的事情,当然除了添加又一个JVM。
- 代理的失败并不会增加困扰。在数据路径上,引入一个代理层并不意味着带来额外开销。客户端不用去关心需要连接的代理。如果因为代理故障而造成的超时,客户端只需要随便选择一个其他的代理。在代理等级不会发生任何分片,它们同样是无状态的。为了扩展吞吐量可以添加更多的代理,代价是额外的成本。代理层只会因为转发被分配资源。集群管理、分片以及集群状态监视都不是代理负责的范围。代理之间并不需要保持一致。
- Twitter中存在实例同时打开10万个连接的情况,服务器也并不会因此发生故障。在Twitter,服务器没理由去关闭一个连接,让它一直打开有助于改善延时。
- 缓存集群被用作后备缓存(look-aside cache)。缓存本身不会负责数据的补充,客户端负责取得一个丢失的键并进行缓存。如果一个节点发生故障,分片会被转移到另一个节点。对故障恢复后的服务器进行同步以保证数据的一致性,所有这些都通过集群管理者完成。在让一个集群更容易理解上,中央观点确实非常重要。
- 测试使用C++来编写代理。C++代理带来了一个显著的性能提升,随后代理层都使用了C和C++。
数据洞察
- 当有调用显示缓存系统失效,而大多数缓存都是正常时,这通常因为客户端的配置错误,或者它们请求的键过多从而导致滥用缓存,当然也可能因为对同一个键多次请求造成服务器或链接饱和。
- 如果通知某个人正在滥用系统,他们很可能会要求出示证据。哪个键?哪个分片的问题?什么样的流量导致了这个问题?因此你需要做足够的度量和分析,从而将证据展示给客户。
- 面向服务的架构并不会给带来问题隔离或者是自动debug,必须对组成系统的每个组件保持足够的能见度。
- 决定在缓存上获得洞察力。缓存使用C编写所以足够快速,因此它可以能其他组件所不能,提供足够的数据,而其他服务不能为每个请求都提供数据。
- 可以实现为每条命令单独建立日志。在10万QPS时,缓存可以记录下所有发生的事情。
- 避免锁和阻塞,特别不能因为磁盘写入而造成阻塞。
- 在每秒100请求和每条日志消息100字节的情况下,每台服务器每秒会记录10MB的数据。当问题发生时,这些数据传输将造成很大的网络开销,大约占10%的带宽,这种开销完全不允许。
- 预计算服务器上的日志以减少开销。设想是已经清楚需要被计算的内容,一个进程负责读取日志,计算并生成一个摘要信息,然后只定期发送主机的摘要信息。对比原始数据,这个体积将微不足道。
- 这些云计算数据通过Storm聚合和存储,并在上面建立可视化系统,你可以基于这些建立容量规划。因为每条日志都可以被捕获,所以你可以做许多事情。
- 对于运营来说,洞察力非常重要。如果出现丢包现象,通常情况下是热键或者是流量峰值导致。
对Redis的希望清单
- 显式的内存管理。
- Deployable(Lua)Scripts。
- 多线程,可以简化集群管理。Twitter有很多高性能服务器,每个主机都拥有100GB以上的内存以及大量的CPU。为了使用一个服务器的所有能力,需要在实体主机商开启许多Redis实例。通过多线程减少需要启动实例数量,从而更容易的进行管理。
学到的知识
- 可预见的规模需求。集群越大、客户越多,你越希望你的服务更可预知和确定。如果只有一个用户和一个问题,你可以迅速的找到并解决这个问题。但是如果你有70个客户,你完全忙不过来。
- 如果你需要在许多分片上推广某个更新,一个分片的速度问题就可能导致全局问题。
- Twitter正在向container环境发展,使用了Mesos作为作业调度器,调度器用来给请求分配CPU、内存等资源数量。当作业占用的资源高于请求时,监视器会直接将它终止。在容器的环境下,Redis会产生一个问题。Redis引入了外部存储碎片,这意味着你要使用更多内存来存储同样的数据。如果不想作业被终止,必须设计一个缓冲的区间。你可能会认为内存碎片率设定在5%就足矣,但是我更愿意多分配10%,甚至是20%的空间作为缓冲。或者在我认为每台主机连接数可能会达到5000时,我将给系统分配支撑1万个连接数的内存,结果会造成很大的浪费。对于当下大多数低延时服务来说,Mesos都不太适合,因此这些作业会与其他作业隔离。
- 在运行时清楚资源使用是非常重要的。在一个大型集群中,总会有问题产生。虽然你一直认为系统很安全,但是事情还是以意料之外的方式发生了。当下很少出现某台机器完全崩溃的情况,比如,在达到10GB的内存上限后,在有空闲内存之前,请求都会被拒绝,造成的后果仅是一小部分请求不能获得自己所需的内存资源,无伤大雅。而垃圾回收问题则是非常麻烦的,它可能降低系统的流量,当下这个问题已经困扰了很多机构。
- 用数据说话。在计算到磁盘和计算到网络之前,查看相对网络速度、CPU速度计磁盘速度是非常有意义的,比如,节点被推送到中央监视服务之前查看的日志综述。除此之外,Redis中的LUA也是给数据提供计算的途径。
- LUA当下还没有在Redis生产环境中实现。响应式脚本意味着服务提供商不能保证他们的SLA,一个被加载的脚本可以做任何事情,因此没有服务提供商会因为添加一些代码铤而走险去破坏SLA。但是对于部署模型来说,意义重大,它将允许代码评审和基准测试,以清晰的计算资源使用和性能。
- Redis作为下一个高性能流处理平台。当下已经拥有pub-sub和scripting,没什么不可以的。
Using Memcached cache from your C# application
Here is the article to read.
ASP.NET internal cache
There is a cache object build into the framework that you could use. The cache object is only ideal for applications hosted in a single server environment.
First of all, an instance of the cache class is created on a per-AppDomain basis and remains valid until that AppDomain is up and running so it’s gone when the app restarts. Furthermore this cache is also not shared between machines in our webfarm meaning all servers still have to retrieve all the data for themselves.
out-of-process and distributed cache
Using an out-of-process and distributed cache then seems ideal since multiple servers and even multiple applications could make use of the same cache as long as they share the key definitions.
ASP.NET internal cache
There is a cache object build into the framework that you could use. The cache object is only ideal for applications hosted in a single server environment.
First of all, an instance of the cache class is created on a per-AppDomain basis and remains valid until that AppDomain is up and running so it’s gone when the app restarts. Furthermore this cache is also not shared between machines in our webfarm meaning all servers still have to retrieve all the data for themselves.
out-of-process and distributed cache
Using an out-of-process and distributed cache then seems ideal since multiple servers and even multiple applications could make use of the same cache as long as they share the key definitions.
Actionable Items
Django’s cache framework
Here is the link.
I like to think about more about the test case, a simple example given by Guo Lisa in her presentation.
I like to think about more about the test case, a simple example given by Guo Lisa in her presentation.
Django: The Web framework for perfectionists with deadlines
Here is the introduction I like to study.
Meet Django
Ridiculously fast
Reassuringly secure
Exceedingly scalable
Why Django?
Ridiculously fast
Fully loaded
Reassuringly secure
Exceedingly scalable
Incredibly versatile
Does Django scale?
Yes. Compared to development time, hardware is cheap, and so Django is designed to take advantage of as much hardware as you can throw at it.
Django uses a "shared-nothing" architecture, which means you can add hardware at any level - database servers, caching servers or Web/application servers.
The framework cleanly separates components such as its database layer and application layer. And it ships with a simple-yet-powerful cache framework.
Meet Django
Ridiculously fast
Reassuringly secure
Exceedingly scalable
Why Django?
Ridiculously fast
Fully loaded
Reassuringly secure
Exceedingly scalable
Incredibly versatile
Does Django scale?
Yes. Compared to development time, hardware is cheap, and so Django is designed to take advantage of as much hardware as you can throw at it.
Django uses a "shared-nothing" architecture, which means you can add hardware at any level - database servers, caching servers or Web/application servers.
The framework cleanly separates components such as its database layer and application layer. And it ships with a simple-yet-powerful cache framework.
F4 - Photo Storage at Facebook
https://www.youtube.com/watch?
F4 - Photo Storage at Facebook
F4 - photo storage at Facebook
world's largest photo storage
BLOB "hotness"
Haystack: 2008
Hot storage design goals
HIgh throughput
- in memory index
- single I/o per request!
- Multiple copies
Failure tolerant
RAIDS
Multiple copies
RAID 6
2 Redudant drivers
2 of 12 = 1.2x replication
Over 3 arrays = 3.6x
RAID 2, RAID 3, RAID 4, RAID 6 explained with diagram
https://www.thegeekstuff.com/
Haystack warm storage
- Redundancy = replcation
Read throughput = replication
Total replication = 3.6x
Warm storage
Redundancy still required
Read throughput the ...
RS Encoding:
Reduundancy
upload request -> web server -> Storage router -> haystack
F4
haystack -> Migration -> F4
Read request -> CDN -> Storage router -> F4
f4: What are we solving
Warm storage problem:
- Need to store (warm) data efficiently
- Storage must be highly fault tolerant
- Read latency should be comparable to haystack
- Load is NOT primary concern
Solution: f4
- 2.x replication factor compared to haystack's 3.6x
- yet more fault tolerant than haystack(!!)
f4: Data splitting RS(5,2)
10G Haystack volume
1G Data blocks
Data blocks
Parity blocks
f4: RS Rebuild
RS decoding
Data blocks parity blocks
f4: Block placement policy
Blocks of each stripe is placed in different racks(=> hosts)
RS(10,4) is used in practice (1.4x)
Tolerant 4 racks(-> 4 disks/hosts) failures
f4 Cell anatomy
f4 storage consists of a set of cells.
One cell resides completely in one data center
Cell consists of 3 kind of nodes: storage, compute, coordinator
The index is distributed across storage needs
f4 Reads
user request -> router -> index read (1), storage nodes, compute, cell
Data read (2)
REads with datacenter failures (2.1X)
Router1, router2, router3
Datacenter1
Datacenter2
Datacenter3
Volume1
volume2
xorVolume
Messaging at Scale at Instagram
Messaging at Scale at Instagram
https://www.youtube.com/watch? v=E708csv4XgY
Chaines tasks
Batch of 10,000 followers per task
Tasks yield successive tasks
much finer-grained load balancing
Failure/Reload penalty low
Other Async tasks
Cross-posting to other networks
search indexing
spam analysis
account deletion
API hook
Gearman framework - load balancer
Gearman in production
persistence horrifically slow, complex
So we ran out of memory and crashed, no recovery
Single core, didn't scale well;
60ms mean submission time for us
Probably should have just used Redis
(Gearman vs Redis)
Celery
Distributed task framework
Highly extensible, pluggable
Mature, Feature rich
Great tooling
Excellent Django support
celeryd
Which broker?
Redis
We already use it
Very fast, eifficient
Polling for task distribution
Messy Non-Synchronous replication
Memory limits task capacity
Beanstalk
Purpose-built task queue
Very fast, efficient
Pushes to Consumers
Spills to disk
No replication
Useless for anything else
RabbitMQ
Reasonablely fast, efficient
Spill-to-disk
Low-maintenance synchronous replication
Excellent celery compatibility
Supports other use cases
We don't know Erlang
Out RabbitMQ Setup
Rabbit MQ 3.0
Clusters of two brokers nodes, Mirrowed
Scale out by adding broker clusters
EC2 c1.xlarge, RAID instance storage
Way overprovisioned
Alerting
We use Sensu
Monitors & alerts on queue length threshold
Uses rabbitmqctl list_queues
Scaling out
Celery only supported 1 broker host last year when we started
Created kombu-multbroker "shim"
Multiple brokers used in a round-robin fashion
Breaks some Celery management tools :(
Concurrency models
multiprocessing (pre-fork)
eventlet
gevent
threads
Problem:
Network-bound tasks sometimes need to take some action
Ruun higer concurrency?
Inefficient :(
Lower batch (prefetch) size?
Min is concurrency count, inefficient :(
Separate slow & fast tasks :)
OUr concurrency levels
fast (14)
feed (12)
default (6)
Problem
Task fails sometimes
Work crashes still lost task
Problem:
Slow tasks monopolize workers
NLP proof gives us choices: to retry or not to retry
problem
Early on, drop task
Publishers confirms
Confirm tasks
Avoid using async tasks as a "backup" mechanism only during failures. It'll probably break.
Better grip on RabbitMQ performance
Utilize result storage
Single cluster for control queues
Eliminate kombu-multibroker
-- study topics --
Sensu
https://docs.sensu.io/sensu- core/1.4/reference/checks/# what-is-a-sensu-check
------------------------------ ----
Study one more topic
AWS Elastic Beanstalk
https://aws.amazon.com/ elasticbeanstalk/
AWS Elastic Beanstalk
Load balancing, provisioning, Application health monitoring, Auto scaling
capacity provisioning, load balancing, auto-scaling, and application health monitoring.
Manageing and configuring servers, databases, load balancers, firewalls, and networks.
https://www.youtube.com/watch?
Chaines tasks
Batch of 10,000 followers per task
Tasks yield successive tasks
much finer-grained load balancing
Failure/Reload penalty low
Other Async tasks
Cross-posting to other networks
search indexing
spam analysis
account deletion
API hook
Gearman framework - load balancer
Gearman in production
persistence horrifically slow, complex
So we ran out of memory and crashed, no recovery
Single core, didn't scale well;
60ms mean submission time for us
Probably should have just used Redis
(Gearman vs Redis)
Celery
Distributed task framework
Highly extensible, pluggable
Mature, Feature rich
Great tooling
Excellent Django support
celeryd
Which broker?
Redis
We already use it
Very fast, eifficient
Polling for task distribution
Messy Non-Synchronous replication
Memory limits task capacity
Beanstalk
Purpose-built task queue
Very fast, efficient
Pushes to Consumers
Spills to disk
No replication
Useless for anything else
RabbitMQ
Reasonablely fast, efficient
Spill-to-disk
Low-maintenance synchronous replication
Excellent celery compatibility
Supports other use cases
We don't know Erlang
Out RabbitMQ Setup
Rabbit MQ 3.0
Clusters of two brokers nodes, Mirrowed
Scale out by adding broker clusters
EC2 c1.xlarge, RAID instance storage
Way overprovisioned
Alerting
We use Sensu
Monitors & alerts on queue length threshold
Uses rabbitmqctl list_queues
Scaling out
Celery only supported 1 broker host last year when we started
Created kombu-multbroker "shim"
Multiple brokers used in a round-robin fashion
Breaks some Celery management tools :(
Concurrency models
multiprocessing (pre-fork)
eventlet
gevent
threads
Problem:
Network-bound tasks sometimes need to take some action
Ruun higer concurrency?
Inefficient :(
Lower batch (prefetch) size?
Min is concurrency count, inefficient :(
Separate slow & fast tasks :)
OUr concurrency levels
fast (14)
feed (12)
default (6)
Problem
Task fails sometimes
Work crashes still lost task
Problem:
Slow tasks monopolize workers
NLP proof gives us choices: to retry or not to retry
problem
Early on, drop task
Publishers confirms
Confirm tasks
Avoid using async tasks as a "backup" mechanism only during failures. It'll probably break.
Better grip on RabbitMQ performance
Utilize result storage
Single cluster for control queues
Eliminate kombu-multibroker
-- study topics --
Sensu
https://docs.sensu.io/sensu-
------------------------------
Study one more topic
AWS Elastic Beanstalk
https://aws.amazon.com/
AWS Elastic Beanstalk
Load balancing, provisioning, Application health monitoring, Auto scaling
capacity provisioning, load balancing, auto-scaling, and application health monitoring.
Manageing and configuring servers, databases, load balancers, firewalls, and networks.
Scaling databases
https://www.youtube.com/watch? v=dkhOZOmV7Fo
Scaling databases
too much load
master - slave1, slave2, replicate
master handles write, slave handles read, a lot of read compared to write.
Downsides
- doesn't increase with speed
- replication log
Too much data in one machine
Master - does not fit into one machine, one computer memory
Too much data
Shard
Hash to multiple machines
1-100 101-200 201-300
downsides
- complex queries
- range query - hit all the machines - merge the result -> sort them in the memory - join becomes difficult
replicate and sharding
One setup for replicate
another setup for sharding
Scaling databases
too much load
master - slave1, slave2, replicate
master handles write, slave handles read, a lot of read compared to write.
Downsides
- doesn't increase with speed
- replication log
Too much data in one machine
Master - does not fit into one machine, one computer memory
Too much data
Shard
Hash to multiple machines
1-100 101-200 201-300
downsides
- complex queries
- range query - hit all the machines - merge the result -> sort them in the memory - join becomes difficult
replicate and sharding
One setup for replicate
another setup for sharding
Wednesday, August 7, 2019
How do I overcome the nervous and anxiety issue?
August 7, 2019
It is my short research. When I am nervous, I tend to look up wechat and then continuously waste a few minutes to get some update. What I really look for is the personal connection. As a single person and also sometime introvert, I like to figure out ways to break those bad habits, and then focus on my own study to prepare most important onsite interview experience.
I just spent a few minutes to watch one of Facebook instagram video provided by Lisa Guo. I quickly understand that I should watch a few videos, every video I will set up a role model for me to follow. So I can calm down and really learn as a senior engineer.
Introduction
It is my short research. When I am nervous, I tend to look up wechat and then continuously waste a few minutes to get some update. What I really look for is the personal connection. As a single person and also sometime introvert, I like to figure out ways to break those bad habits, and then focus on my own study to prepare most important onsite interview experience.
One idea
I just spent a few minutes to watch one of Facebook instagram video provided by Lisa Guo. I quickly understand that I should watch a few videos, every video I will set up a role model for me to follow. So I can calm down and really learn as a senior engineer.
How to make good learning experience during preparation?
August 7, 2019
It is my curiosity research. I like to learn as many topics related to Amazon and Facebook during my system design and product design preparation. If I cannot understand the topics, I can quickly go back to those system design tutorial video, and learn the basics first.
I like to learn a few areas related to Amazon AWS.
Introduction
It is my curiosity research. I like to learn as many topics related to Amazon and Facebook during my system design and product design preparation. If I cannot understand the topics, I can quickly go back to those system design tutorial video, and learn the basics first.
Amazon AWS
I like to learn a few areas related to Amazon AWS.
Vacation from August 12 to August 20, 2019
August 7, 2019
Introduction
It is the important for me to prepare and schedule a vacation from August 12 to August 20, 2019. I like to spend time to learn system design, and also work on algorithm and data structure during the break.
Introduction
It is the important for me to prepare and schedule a vacation from August 12 to August 20, 2019. I like to spend time to learn system design, and also work on algorithm and data structure during the break.
Subscribe to:
Posts (Atom)



