Wednesday, August 12, 2020

System design: Facebook TAO data store

 August 12, 2020

Introduction

It is the third day I like to learn Facebook TAO data store. I came cross an ex-Microsoft researcher and his study notes. 

Notes study

I just quickly copy the notes into my blog, and then read carefully. I like to watch the presentation again to start my study day. 

Here is the link. 


TAO

  • Geographically distributed, read-optimized, graph data store.
  • Favors availability and efficiency over consistency.
  • Developed by and used within Facebook (social graph).
  • Link to paper.

Before TAO

  • Facebook's servers directly accessed MySQL to read/write the social graph.
  • Memcache used as a look-aside cache.
  • Had several issue:
    • Inefficient edge list - A key-value store is not a good design for storing a list of edges.
    • Distributed Control Logic - In look-aside cache architecture, the control logic runs on the clients which increase the number of failure modes.
    • Expensive Read-After-Write Consistency - Facebook used asynchronous master-slave replication for MySQL which introduced a time lag before latest data would reflect in the local replicas.

TAO Data Model

  • Objects

    • Typed nodes (type is denoted by otype)
    • Identified by 64-bit integers.
    • Contains data in the form of key-value pairs.
    • Models users and repeatable actions (eg comments).
  • Associations

    • Typed directed edges between objects (type is denoted by atype)
    • Identified by source object id1, atype and destination object id2.
    • Contains data in the form of key-value pairs.
    • Contains a 32-bit time field.
    • Models actions that happen at most once or records state transition (eg like)
    • Often inverse association is also meaningful (eg like and liked by).

Query API

  • Support to create, retrieve, update and delete objects and associations.
  • Support to get all associations (assoc_get) or their count(assoc_count) based on starting node, time, index and limit parameters.

TAO Architecture

Storage Layer

  • Objects and associations stored in MySQL.
  • TAO API mapped to SQL queries.
  • Data divided into logical shards.
  • Objects bound to the shard for their lifetime(shard_id is embedded in id).
  • Associations stored on the shard of its id (for faster association query).

Caching Layer

  • Consists of multiple cache servers (together form a tier).
  • In memory, LRU cache stores objects, association lists, and association counts.
  • Write operation on association list with inverse involves writing 2 shards (for id1 and id2).
  • The client sends the query to cache layer which issues inverse write query to shard2 and once that is completed, a write query is made to shard1.
  • Write failure leads to hanging associations which are repaired by an asynchronous job.

Leaders and Followers

  • A single, large tier is prone to hot spots and square growth in terms of all-to-all connections.
  • Cache split into 2 levels - one leader tier and multiple follower tiers.
  • Clients communicate only with the followers.
  • In the case of read miss/write, followers forward the request to the leader which connects to the storage layer.
  • Eventual consistency maintained by serving cache maintenance messages from leaders to followers.
  • Object update in leaders leads results in invalidation message to followers.
  • Leader sends refill message to notify about association write.
  • Leaders also serialize concurrent writes and mediates thundering herds.

Scaling Geographically

  • Since workload is read intensive, read misses are serviced locally at the expense of data freshness.
  • In the multi-region configuration, there are master-slave regions for each shard and each region has its own followers, leader, and database.
  • Database in the local region is a replica of the database in the master region.
  • In the case of read miss, the leader always queries the region database (irrespective of it being the master database or slave database).
  • In the case of write, the leader in the local region would forward the request to database in the master region.

Optimisations

Caching Server

  • RAM is partitioned into arena to extend the lifetime of important data types.
  • For small, fixed-size items (eg association count), a direct 8-way associative cache is maintained to avoid the use of pointers.
  • Each atype is mapped to 16-bit value to reduce memory footprint.

Cache Sharding and Hot Spots

  • Load is balanced among followers through shard cloning(reads to a shard are served by multiple followers in a tier).
  • Response to query include the object's access rate and version number. If the access rate is too high, the object is cached by the client itself. Next time when the query comes, the data is omitted in the reply if it has not changed since the previous version.

High Degree Objects

  • In the case of assoc_count, the edge direction is chosen on the basis of which node (source or destination) has a lower degree (to optimize reading the association list).
  • For assoc_get query, only those associations are searched where time > object's creation time.

Failure Detection and Handling

  • Aggressive network timeouts to detect (potential) failed nodes.

Database Failure

  • In the case of master failure, one of the slaves take over automatically.
  • In case of slave failure, cache miss are redirected to TAO leader in the region hosting the database master.

Leader Failure

  • When a leader cache server fails, followers route read miss directly to the database and write to a replacement leader (chosen randomly from the leader tier).

Refill and Invalidation Failures

  • Refill and invalidation are sent asynchronously.
  • If the follower is not available, it is stored in leader's disk.
  • These messages will be lost in case of leader failure.
  • To maintain consistency, all the shards mapping to a failed leader are invalidated.

Follower Failure

  • Each TAO client is configured with a primary and backup follower tier.
  • In normal mode, the request is made to primary tier and in the case of its failure, requests go to backup tier.
  • Read after write consistency may be violated if failing over between different tiers (read reaches the failover target before writer's refill or invalidate).

Actionable items

I should compare myself to a professional researcher, and then I should write down my own notes to compare. It is tough for me to learn in such a short time. 


Boxlight Inc: Linkedin profile

 I like to spend 10 minutes to work on research on boxlight inc. 


Stock research: One hour 30 minutes

 August 12, 2020

Introduction

It is a normal day, my third vacation day. I spent morning time to work on stock research. However, my research is not so efficient, I spent time to read news on fenviz.com, and then I did visit a few other websites. I like to write down this research. 

8:00 AM - 9:29 AM

I spent time to learn basica, read finance news, and also think about my portfolios. I just could not believe that it is such pleasant experience to study market. 

My friend on wechat group shared her purchase of BOXL stock, and her profit with over 24% gains in one day. 

I am not sharp on this stock market research. I should expand my search on education stock, and also penny stocks as well. 



Why Hudbay Minerals Stock Is Soaring Today

 Here is the article. 

What happened

Shares of Hudbay Minerals (NYSE:HBM) had surged more than 10% by 10:45 a.m. EDT on Wednesday. Driving up the mining stock was its stronger-than-expected second-quarter results.  

So what

Hudbay Minerals reported an adjusted loss of $0.15 per share, a $0.05 per share improvement from what analysts expected. Several factors drove the narrower loss, including strong production and cost performance at its operations in Manitoba, and rising production of all other metals elsewhere, including record gold output. That strong gold production came at an ideal time as prices soared during the period. 

The mining company's excellent performance in Manitoba during the first half of the year has it on track to achieve its full-year production and cost guidance for those operations despite the impact of COVID-19. But the company did warn that the pandemic had a significant effect on its operations in Peru, where the government required it to suspend work in mid-March. While the company resumed operations at its Constancia mill in mid-May and returned to normal levels by early July, it had to cut its 2020 production forecast for Peru. That's because of the Constancia suspension and pushing back the start of mining at Pampacancha until early next year. 

Boxl: Boxlight and Samsung Enter into Strategic Partnership for the US Education Market

 Here is the article. 

My notes from study:

  1. Boxlight’s OKTOPUS software
  2. OKTOPUS Blend, a library of over 10,000 master-teacher created lessons, and GameZones, a gaming component that features over 90 games that focus on concepts and skills appropriate for K-8 students
  3. Boxlight’s professional development courses

Jointly offer bundled education solutions to include displays, classroom software and professional development.

Boxlight Corporation (NASDAQ:BOXL), a leading provider of interactive technology solutions for the global education market, and Samsung Electronics America today announced a strategic partnership to bundle classroom displays, classroom software and professional development for the education market.

The education bundles combine Samsung display innovations with Boxlight’s OKTOPUS education software and online professional development courses. The bundles provide teachers a complete solution that can be easily adopted to create an engaging learning environment for students.

Samsung display solutions for education include digital signage, interactive Flip 2 and QBH-TR displays. Boxlight’s OKTOPUS software provides a presentation solution that offers subject specific tools for teaching, assessment and polling functions, grading and reporting features, and personal device collaboration. OKTOPUS also allows existing third-party lessons to be uploaded and used in the teaching software. In addition to OKTOPUS, the education bundles include OKTOPUS Blend, a library of over 10,000 master-teacher created lessons, and GameZones, a gaming component that features over 90 games that focus on concepts and skills appropriate for K-8 students.

Boxlight’s professional development courses provide training for both front of classroom Samsung displays and the OKTOPUS education software. The courses are offered online and are interactive, requiring learners to apply new skills via practice exercises. Teachers can check their knowledge and review course content as needed.

"Schools are constantly looking for ways to better collaborate and Samsung is investing resources into the market to make new education technology possible," said Chris Mertens, Vice President Sales, Display Division, Samsung Electronics America. "Through our partnership with Boxlight, Samsung is proud to usher in the next generation of education tools. Samsung and Boxlight solutions offer educators endless productivity possibilities by condensing multiple tools and innovations required for lessons as part of the education bundles that will be available later this year."

Actionable Items

My friend bought 4000 share on August 11, 2020, $1.5 a share, and sold them $2.5/ share. 

Active investor: How to find investment opportunities?

 August 12, 2020

Introduction

There are so many things for me to learn since I am a beginner as an active investor. I do like to take some time to learn more about business, and then wait for my next project. I like to write down a few ways for me to look into stock market. 

How to find investment opportunities? 

I do think that I was too short-sighted. I purchased 300 shares of MGM in June, 2020, but I only held them less than one week. I made profit around $150.00 dollars. But I could have hold them until August 12, I should have made gains around $1500 dollars. 

So I have to learn how to invest in a longer term. Do not worry about market value. Try to think calmly and stay positive. 

It is not easy for me to stay focus on holding stocks with distressed value since coronavirus. I should learn how MGM stays in business and handle the coronavirus issues. 



Tuesday, August 11, 2020

Mozilla lays off 250

 Here is the article. 


Mozilla  today announced a major restructuring of its commercial arm, the Mozilla Corporation, that will see about 250 employees lose their jobs and the shuttering of the organization’s operations in Taipei, Taiwan. This move comes after the organization already laid off about 70 employees earlier this year. The most recent numbers from 2018 put Mozilla at about 1,000 employees worldwide.

Citing falling revenues because of the global pandemic, Mozilla’s executive chairwoman and CEO Mitchell Baker  said in an internal message that the company’s pre-COVID plans were no longer feasible.

“Pre-COVID, our plan for 2020 was a year of change: building a better internet by accelerating product value in Firefox, increasing innovation, and adjusting our finances to ensure financial stability over the long term,” Baker writes. “We started with immediate cost-saving measures such as pausing our hiring, reducing our wellness stipend and cancelling our All-Hands. But COVID-19 has accelerated the need and magnified the depth for these changes. Our pre-COVID plan is no longer workable. We have talked about the need for change — including the likelihood of layoffs — since the spring. Today these changes become real.”

System design: Common system design interview questions

 Here is the link. 

Design Twitter timeline and search ( or Facebook feed and search)

Design Mint.com

Design the data structures for a social network

Design a key-value store for a search engine

Design Amazon's sales ranking by category feature

Design a system that scales to millions of users on AWS


System design: Designing a URL shortener

 Here is the article. 

My study notes:

I just could not believe that the author puts together the comparison what is good, not good. I am so surprised to learn that I am not so smart to figure out those things. A lot of ideas are new to me. 

Here are highlights:

  1. Total number of unique domains - 10s of thousands a day, 10,000, 000/ day or 100/sec;
  2. 10B/ day, 100,000/sec redirect requests
  3. Peak time vs average request per second
  4. 5 minutes lag time aggregated data
Encode URL - ideas to work on it, make it unique, how many bits? 

Ask about the lifespan of the aliases and design a system that purges aliases past their expiration.

So using Key value store instead of relational database - 

zooKeeper - 

The content from the above article. 

Design a system to take user-provided URLs and transform them to a shortened URLs that redirect back to the original. Describe how the system works. How would you allocate the shorthand URLs? How would you store the shorthand to original URL mapping? How would you implement the redirect servers? How would you store the click stats?

Assumptions: I generally don’t include these assumptions in the initial problem presentation. Good candidates will ask about scale when coming up with a design.

  • Total number of unique domains registering redirect URLs is on the order of 10s of thousands
  • New URL registrations are on the order of 10,000,000/day (100/sec)
  • Redirect requests are on the order of 10B/day (100,000/sec)
  • Remind candidates that those are average numbers - during peak traffic (either driven by time, such as ‘as people come home from work’ or by outside events, such as ‘during the Superbowl’) they may be much higher.
  • Recent stats (within the current day) should be aggregated and available with a 5 minute lag time
  • Long look-back stats can be computed daily

Assumptions 

1B new URLs per day, 100B entries in total the shorter, the better show statics (real-time and daily/monthly/yearly)

Encode Url 

http://blog.codinghorror.com/url-shortening-hashes-in-practice/

Choice 1. md5(128 bit, 16 hex numbers, collision, birthday paradox, 2^(n/2) = 2^64) truncate? (64bit, 8 hex number, collision 2^32), Base64.

  • Pros: hashing is simple and horizontally scalable.
  • Cons: too long, how to purify expired URLs?

Choice 2. Distributed Seq Id Generator. (Base62: a~z, A~Z, 0~9, 62 chars, 62^7), sharding: each node maintains a section of ids.

  • Pros: easy to outdate expired entries, shorter
  • Cons: coordination (zookeeper)

KV store 

MySQL(10k qps, slow, no relation), KV (100k qps, Redis, Memcached)

A great candidate will ask about the lifespan of the aliases and design a system that purges aliases past their expiration.

Followup 

Q: How will shortened URLs be generated?

  • A poor candidate will propose a solution that uses a single id generator (single point of failure) or a solution that requires coordination among id generator servers on every request. For example, a single database server using an auto-increment primary key.
  • An acceptable candidate will propose a solution using an md5 of the URL, or some form of UUID generator that can be done independently on any node. While this allows distributed generation of non- colliding IDs, it yields large “shortened” URLs
  • A good candidate will design a solution that utilizes a cluster of id generators that reserve chunks of the id space from a central coordinator (e.g. ZooKeeper) and independently allocate IDs from their chunk, refreshing as necessary.

Q: How to store the mappings?

  • A poor candidate will suggest a monolithic database. There are no relational aspects to this store. It is a pure key-value store.
  • A good candidate will propose using any light-weight, distributed store. MongoDB/HBase/Voldemort/etc.
  • A great candidate will ask about the lifespan of the aliases and design a system that purges aliases past their expiration

Q: How to implement the redirect servers?

  • A poor candidate will start designing something from scratch to solve an already solved problem
  • A good candidate will propose using an off-the-shelf HTTP server with a plug-in that parses the shortened URL key, looks the alias up in the DB, updates click stats and returns a 303 back to the original URL. Apache/Jetty/Netty/tomcat/etc. are all fine.

Q: How are click stats stored?

  • A poor candidate will suggest write-back to a data store on every click
  • A good candidate will suggest some form of aggregation tier that accepts clickstream data, aggregates it, and writes back a persistent data store periodically

Q: How will the aggregation tier be partitioned?

  • A great candidate will suggest a low-latency messaging system to buffer the click data and transfer it to the aggregation tier.
  • A candidate may ask how often the stats need to be updated. If daily, storing in HDFS and running map/reduce jobs to compute stats is a reasonable approach If near real-time, the aggregation logic should compute stats

Q: How to prevent visiting restricted sites?

  • A good candidate can answer with maintaining a blacklist of hostnames in a KV store.
  • A great candidate may propose some advanced scaling techniques like bloom filter.

System design: INTRODUCTION Everything you need to know about Kafka in 10 minutes

 Here is the link. 

My notes:

Log - replicate a few times - think like event 

monolithic - think about small ones, -> log, build new service -> gauge, consuming messages - batch process overnight -> build service 

topic, event, log, monolith - 

distributed log things - tend to build on system - start to use kafka - logs, 

What is event streaming?

Event streaming is the digital equivalent of the human body's central nervous system. It is the technological foundation for the 'always-on' world where businesses are increasingly software-defined and automated, and where the user of software is more software.

Technically speaking, event streaming is the practice of capturing data in real-time from event sources like databases, sensors, mobile devices, cloud services, and software applications in the form of streams of events; storing these event streams durably for later retrieval; manipulating, processing, and reacting to the event streams in real-time as well as retrospectively; and routing the event streams to different destination technologies as needed. Event streaming thus ensures a continuous flow and interpretation of data so that the right information is at the right place, at the right time.

What can I use event streaming for?

Event streaming is applied to a wide variety of use cases across a plethora of industries and organizations. Its many examples include:

  • To process payments and financial transactions in real-time, such as in stock exchanges, banks, and insurances.
  • To track and monitor cars, trucks, fleets, and shipments in real-time, such as in logistics and the automotive industry.
  • To continuously capture and analyze sensor data from IoT devices or other equipment, such as in factories and wind parks.
  • To collect and immediately react to customer interactions and orders, such as in retail, the hotel and travel industry, and mobile applications.
  • To monitor patients in hospital care and predict changes in condition to ensure timely treatment in emergencies.
  • To connect, store, and make available data produced by different divisions of a company.
  • To serve as the foundation for data platforms, event-driven architectures, and microservices.

Apache Kafka® is an event streaming platform. What does that mean?

Kafka combines three key capabilities so you can implement your use cases for event streaming end-to-end with a single battle-tested solution:

  1. To publish (write) and subscribe to (read) streams of events, including continuous import/export of your data from other systems.
  2. To store streams of events durably and reliably for as long as you want.
  3. To process streams of events as they occur or retrospectively.

And all this functionality is provided in a distributed, highly scalable, elastic, fault-tolerant, and secure manner. Kafka can be deployed on bare-metal hardware, virtual machines, and containers, and on-premises as well as in the cloud. You can choose between self-managing your Kafka environments and using fully managed services offered by a variety of vendors.

How does Kafka work in a nutshell?

Kafka is a distributed system consisting of servers and clients that communicate via a high-performance TCP network protocol. It can be deployed on bare-metal hardware, virtual machines, and containers in on-premise as well as cloud environments.

Servers: Kafka is run as a cluster of one or more servers that can span multiple datacenters or cloud regions. Some of these servers form the storage layer, called the brokers. Other servers run Kafka Connect to continuously import and export data as event streams to integrate Kafka with your existing systems such as relational databases as well as other Kafka clusters. To let you implement mission-critical use cases, a Kafka cluster is highly scalable and fault-tolerant: if any of its servers fails, the other servers will take over their work to ensure continuous operations without any data loss.

Clients: They allow you to write distributed applications and microservices that read, write, and process streams of events in parallel, at scale, and in a fault-tolerant manner even in the case of network problems or machine failures. Kafka ships with some such clients included, which are augmented by dozens of clients provided by the Kafka community: clients are available for Java and Scala including the higher-level Kafka Streams library, for Go, Python, C/C++, and many other programming languages as well as REST APIs.

Main Concepts and Terminology

An event records the fact that "something happened" in the world or in your business. It is also called record or message in the documentation. When you read or write data to Kafka, you do this in the form of events. Conceptually, an event has a key, value, timestamp, and optional metadata headers. Here's an example event:

  • Event key: "Alice"
  • Event value: "Made a payment of $200 to Bob"
  • Event timestamp: "Jun. 25, 2020 at 2:06 p.m."

Producers are those client applications that publish (write) events to Kafka, and consumers are those that subscribe to (read and process) these events. In Kafka, producers and consumers are fully decoupled and agnostic of each other, which is a key design element to achieve the high scalability that Kafka is known for. For example, producers never need to wait for consumers. Kafka provides various guarantees such as the ability to process events exactly-once.

Events are organized and durably stored in topics. Very simplified, a topic is similar to a folder in a filesystem, and the events are the files in that folder. An example topic name could be "payments". Topics in Kafka are always multi-producer and multi-subscriber: a topic can have zero, one, or many producers that write events to it, as well as zero, one, or many consumers that subscribe to these events. Events in a topic can be read as often as needed—unlike traditional messaging systems, events are not deleted after consumption. Instead, you define for how long Kafka should retain your events through a per-topic configuration setting, after which old events will be discarded. Kafka's performance is effectively constant with respect to data size, so storing data for a long time is perfectly fine.

Topics are partitioned, meaning a topic is spread over a number of "buckets" located on different Kafka brokers. This distributed placement of your data is very important for scalability because it allows client applications to both read and write the data from/to many brokers at the same time. When a new event is published to a topic, it is actually appended to one of the topic's partitions. Events with the same event key (e.g., a customer or vehicle ID) are written to the same partition, and Kafka guarantees that any consumer of a given topic-partition will always read that partition's events in exactly the same order as they were written.

Figure: This example topic has four partitions P1–P4. Two different producer clients are publishing, independently from each other, new events to the topic by writing events over the network to the topic's partitions. Events with the same key (denoted by their color in the figure) are written to the same partition. Note that both producers can write to the same partition if appropriate.

To make your data fault-tolerant and highly-available, every topic can be replicated, even across geo-regions or datacenters, so that there are always multiple brokers that have a copy of the data just in case things go wrong, you want to do maintenance on the brokers, and so on. A common production setting is a replication factor of 3, i.e., there will always be three copies of your data. This replication is performed at the level of topic-partitions.

This primer should be sufficient for an introduction. The Design section of the documentation explains Kafka's various concepts in full detail, if you are interested.

Kafka APIs

In addition to command line tooling for management and administration tasks, Kafka has five core APIs for Java and Scala:

  • The Admin API to manage and inspect topics, brokers, and other Kafka objects.
  • The Producer API to publish (write) a stream of events to one or more Kafka topics.
  • The Consumer API to subscribe to (read) one or more topics and to process the stream of events produced to them.
  • The Kafka Streams API to implement stream processing applications and microservices. It provides higher-level functions to process event streams, including transformations, stateful operations like aggregations and joins, windowing, processing based on event-time, and more. Input is read from one or more topics in order to generate output to one or more topics, effectively transforming the input streams to output streams.
  • The Kafka Connect API to build and run reusable data import/export connectors that consume (read) or produce (write) streams of events from and to external systems and applications so they can integrate with Kafka. For example, a connector to a relational database like PostgreSQL might capture every change to a set of tables. However, in practice, you typically don't need to implement your own connectors because the Kafka community already provides hundreds of ready-to-use connectors.

System design: What is Apache Kafka?

 Here is the article. 

What I like to learn from the article. Karfka - it was my study topic back in 2019. I do think that Karfka is a good idea using logging to allow distributed processing, kind of liking ACID properties. 

From the above article. 

Why use Apache Kafka? 

Its abstraction is a queue and it features

  • a distributed pub-sub messaging system that resolves N^2 relationships to N. Publishers and subscribers can operate at their own rates.
  • super fast with zero-copy technology
  • support fault-tolerant data persistence

It can be applied to

  • logging by topics
  • messaging system
  • geo-replication
  • stream processing

Why is Kafka so fast? 

Kafka is using zero copy in which that CPU does not perform the task of copying data from one memory area to another.

Without zero copy:

  1. Producer publishes messages to a specific topic.
    • Write to in-memory buffer first and flush to disk.
    • append-only sequence write for fast write.
    • Available to read after write to disks.
  2. Consumer pulls messages from a specific topic.
    • use an “offset pointer” (offset as seqId) to track/control its only read progress.
  3. A topic consists of partitions, load balance, partition (= ordered + immutable seq of msg that is continually appended to)
    • Partitions determine max consumer (group) parallelism. One consumer can read from only one partition at the same time.

How to serialize data? Avro

What is its network protocol? TCP

What is a partition’s storage layout? O(1) disk read

How to tolerate fault? 

In-sync replica (ISR) protocol. It tolerates (numReplicas - 1) dead brokers. Every partition has one leader and one or more followers.

Total Replicas = ISRs + out-of-sync replicas

  1. ISR is the set of replicas that are alive and have fully caught up with the leader (note that the leader is always in ISR).
  2. When a new message is published, the leader waits until it reaches all replicas in the ISR before committing the message.
  3. If a follower replica fails, it will be dropped out of the ISR and the leader then continues to commit new messages with fewer replicas in the ISR. Notice that now, the system is running in an under replicated mode. If a leader fails, an ISR is picked to be a new leader.
  4. Out-of-sync replica keeps pulling message from the leader. Once catches up with the leader, it will be added back to the ISR.

Is Kafka an AP or CP system in CAP theorem? 

Jun Rao says it is CA, because “Our goal was to support replication in a Kafka cluster within a single datacenter, where network partitioning is rare, so our design focuses on maintaining highly available and strongly consistent replicas.”

However, it actually depends on the configuration.

  1. Out of the box with default config (min.insync.replicas=1, default.replication.factor=1) you are getting AP system (at-most-once).

  2. If you want to achieve CP, you may set min.insync.replicas=2 and topic replication factor of 3 - then producing a message with acks=all will guarantee CP setup (at-least-once), but (as expected) will block in cases when not enough replicas (<2) are available for particular topic/partition pair.


System design: How we efficiently implemented consistent hashing

 Here is the article. 



System design: Consistent hashing

 Here is the article. 

It is interesting to learn how to explain consistent hashing in his writing. I may like to add some review on his writing as well. 


From the article. 

Consistent hashing, and was first described by Karger et al. at MIT in an academic paper from 1997 (according to Wikipedia).

Consistent Hashing is a distributed hashing scheme that operates independently of the number of servers or objects in a distributed hash table by assigning them a position on an abstract circle, or hash ring. This allows servers and objects to scale without affecting the overall system.

Distributed hash key value:

To ensure object keys are evenly distributed among servers, we need to apply a simple trick: To assign not one, but many labels (angles) to each server. So instead of having labels A, B and C, we could have, say, A0 .. A9, B0 .. B9 and C0 .. C9, all interspersed along the circle. The factor by which to increase the number of labels (server keys), known as weight, depends on the situation (and may even be different for each server) to adjust the probability of keys ending up on each. For example, if server B were twice as powerful as the rest, it could be assigned twice as many labels, and as a result, it would end up holding twice as many objects (on average).

Removing a server results in its object keys being randomly reassigned to the rest of the servers, leaving all other keys untouched:

Example:

Can you explain how you got 58.8

KEY HASH ANGLE (DEG)
"john" 1633428562 58.8
"bill" 7594634739 273.4
"jane" 5000799124 180
"steve" 9787173343 352.3
"kate" 3421657995 123.2
"A" 5572014558 200.6
"B" 8077113362 290.8
"C" 2269549488 81.7


What I understood was, you take the ratio of the key with the max num, and multiply it by 360(total angle in the circle).
So, max in our case is 10^10(he mentioned it).
So, for John, angle = (1633428562 / 10^10) * 360 = 58.8


Exactly. Just keep in mind we don't really need angles in the actual implementation (the circle is just a way to visualize it).


Netflix tech talk: Netflix architecture evolution

 Here is the link. 

Architecture evolution

Sessions - Oracle database - scale up - no scale out - ad hoc ... painful, not evolved 


Real time data - gen 1 pain points

- scalability - DB scaled up not out

- Event data analytics - ad hoc

- Fixed schema

More expensive Oracle hardware - computer, cannot scale out

Real time data - gen 2 motivation

Scalability - scale out not up 

Flexible schema - key/value attributes

Service oriented

Real time data - gen 2 pain points

- scale out - resharding was painful

- performance - hot spots

- Disaster recovery - simpleDB had no backups

Real time data - gen 3 landscape

- Cassandra 0.6

- Before SSDs in AWS

- Netflix in 1 AWS region

Real time data - gen 3 motivations

- order of magnitude increase in requests

- scalability - actually scale out rather than up 

Real time data - gen 3 writes

start - stop 

viewing service - 

gen 3 - Cluster Scale 

cluster - scale 

Real time Data - gen 3 pain points

- stateful tier - hot spots, multi-region complexity

- monolithic service

-read-modify-write poorly suited for memcached


Real Time Data - gen 3 learnings

- Distributed stateful systems are hard - go stateless, use C*/ memcached/redis...

- Decompose into microservices


Real Time Data - gen 4

stateless Microservices

- stream state/ event collectors

- data processors

- data services

- data feeds


Session analytics

- summarize detailed event data

- non-real time, but near real time

- some shared logic with real time


Session analytics - processing 

Storage - processing 

processing - storm - mantis, Samza 


Polygot persistence - one size fits all doesn't fit all

Strong opinions, loosely held - design for long term, but be open to redesigns



Viewing service - 50 data partitions 

Scale out - resharding was painful 

Performance - hot spots 

Disaster recovery 

NoSQL - > MemCache, Cassandra - 

gen 3 motivation

order of magnitude - include ... 

Write / read stateful tier - active sessions, latest positions, View summary - > sanpshot, viewing history Memcached 

Access - ...