Monday, August 12, 2019

Case study: system design - instagram

August 12, 2019

Introduction


It is very good learning experience for me to learn how to design instagram based on grokking system design content. I like to do a case study and then allow myself to learn how to work on design with very good structure, gathering requirement, high level design, and dive deep, and tradeoff and quantitative analysis.

Case study


Here are the highlights I should remember to be evaluated in system design.

1. Problem exploration
2. Design approach
3. Data management
4. Tradeoffs
5. Deep dive
6. Quantitative analysis

I have to push myself to learn and how to design a simpler version of instagram. 

Functional requirements

1. User should be able to upload/download/view photos
2. User can perform searches based on photo/video titles
3. User can follow other users
4. The system should be able to generate and display a user's News Feed consisting of top photos from all the people the user follows. 

Non-functional requirements

1. Our service needs to be highly available.
2. The acceptable latency of the system is 200ms for News Feed generation. 
3. Consistency can take a hit (in the interest of availability), if a user doesn't see a photo for a while; it should be find. 
4. The system should be highly reliable; any uploaded photo or video should never be lost. 

Not in scope: 
Adding tags to photos, searching photos on tags, commenting on photos, tagging users to photos, who to follow, etc. 

So, three things:

Functional requirements
Non-functional requirements
Not in scope 


Design approach



The system would be read-heavy, so we will focus on building a system that can retrieve photos quickly. 

1. Management of storage is important
2. Low latency is expected while viewing photos
3. Data should be 100% reliable. If a user uploads a photo, the system will guarantee that it will never be lost. 

- ideas related to management of storage: distributed file storage like HDFS or S3. 



4. Capacity estimation and constraints 

Let us assume we have 500M total users, with 1M daily active users. 
2M new photos every day, 23 new photos every second. 
Average photo file size => 200KB
total space required for 1 day of photos
   2M * 200KB => 400 GB 
total space required for 10 years:
  400GB * 365 ( days a year) * 10 (years) ~= 1425 TB

5. Database schema - Design database tables first 

It is a good idea to design the database schema. The database schema would help to understand the data flow among various components and later would guide towards data partitioning. 

It is important to write down tables with columns in the table, size and data type, primary key. 

Tables:
Photo 
User
UserFollow

We can consider using an RDBMS like MySQL, there is an issue to scale them. 
We can consider using a distributed key-value store instead of using NoSQL. 

I need to learn how to dive deep and do quantitative analysis. 

It is important for me to write down table, column name, data type, size, so I can calculate in detail how many bytes for each row. 

Let me skip the calculation. But I have to look into the size of table, 3.7 TB for all tables. 

Photo: 
284 bytes - each row in photo's table
one day, 2 million new photos get uploaded every day, 0.5 GB of storage for one day. 
For 10 years 1.88 TB of storage. 

I need to think carefully about Sharding. 

Assume that a web server can have a maximum of 500 connections at any time. 
Upload photos can be a bottleneck. 
The idea is to separate reads from writes into two different service. We will have dedicated servers for reads and different servers for writes to ensure that uploads don't affect the system. 



Data sharing 

I like to learn how to figure out different schemes for metadata sharing:

1. Partitioning based on UserID 

If one DB shard is 1TB, we will need four shards to store 3.7 TB of data. Let's assume that we keep 10 shards. 

To uniquely identify any photo in the system, it is a working idea to append shard number with each PhotoID. 

How can we generate PhotoIDs?
What are the different issues with this partitioning scheme? 

It is challenging for me to think about the issues with the partitioning scheme. 

1. How are hot users handled? 
2. Some users will have a lot of photos compared to others, non-uniform distribution of storage
3. One user's photos are on multiple shards, will it cause higher latencies?
4. One shard for a user's all photos, if the shard is down or high latency if it is serving high load. 

b. Partitioning based on PhotoID

1. wouldn't this key generating DB be a single point of failure? 
Yes. It would be. A workaround for that is to use two databases with one generating even number and the other odd numbered. 


Ranking and News Feed Generation



I am so surprised to read the content. I do not remember at all I read those content in 2018 when I prepare for Amazon onsite on June 6, 2018. I just could not believe that I did not write down any note on this topic at all. 


News feed creation with sharded data



Given the task to create the News Feed for any given user is to fetch the latest photos from all people the user follows. 

How to come out a mechanism to sort photos on their time of creation. One of ideas is to add photo creation time part of the PhotoID. As the primary index is available on PhotoID, it will be easy. 


Follow up 


Nov. 15, 2019

It is the first time I like to write down my review after six three months. It is true that I have no idea how good or bad I will be in terms of distributed system design. But I like to have some structure in terms of learning. Read Grokking system design, take notes, and then review notes once a while. 


It is better to spend 20 minutes to read and think about what to learn from this blog. 

Dec. 11, 2019
I spent time to review the blog, I have onsite from dialpad on Dec. 16, four days away. 

Comparing Hadoop to Distributed Databases” on page 414.

I like to read the topic “Comparing Hadoop to Distributed Databases” on page 414.

Stocks slammed as rates tumble

Here is the link.

10-YR T-NOTE




Course Development: Behind the Scenes

Here is the link.


Study Habits from the Udacity Data Team

Here is the link.


10 reasons I like Zumba in Canada place

August 12, 2019

Introduction


It is my personal finance research. I like to have frugal life style. How can I do that and also enjoy the life style? I talked to my coworkers and they laughed about it. I told them that I will start to work on 401 K and TFSA portfolio, and then in two months I showed the portfolio starting from June 2019. I think that they like me, and I was invited to attend Zumba dance in Canada place.

10 reasons I like Zumba in Canada place


Here is the link for the event.

It is hard for me to be honest. My coworkers used to ask me how old I am. And also I was surprised to learn that I have to control my weight again as well. I think that being frugal is not too difficult. I just need to attend more events in the city, and then get connected to more people.

1. I like Zumba dance. It can be a good warmup exercise for my tennis sports;
2. I like the Canada place, and weather is so nice and waterfront is such a great place to visit;
3. I need to take some time off and do more workout.
4. ...



The schedule I choose first day vacation

August 12, 2019

Introduction


It is a tough task to put schedule together for me to study and also keep myself active in the same time. I just started my first day vacation today.


My schedule


I woke up around 6:30 AM, so I played a few videos related to Udacity and got up 7:30 AM.

First thing, I took a walk around my neighborhood 20 - 30 minutes.
Study the book - Design large data intensive application
Lunch - 20 - 30 minute lunch
Walk around neighborhood 1:30 - 1:50 PM - 20 minutes break
Continue to study
4:40 PM - 5:30 PM Go to waterfront to attend Zumba class
5:30 PM - 6:30 PM Zumba class, here is the link.
6:30 PM - 7:20 PM commute to home
Continue to study
8:30 - 9:00 20 minutes walk
9:00 - 12:00 PM 3 hours study nonstop


How we built a data pipeline with Lambda Architecture using Spark/Spark Streaming

Here is the article to read written by WarmartLab senior engineer.


Lambda architecture

Here is the link.



Lambda architecture is a data-processing architecture designed to handle massive quantities of data by taking advantage of both batch and stream-processing methods. This approach to architecture attempts to balance latency, throughput, and fault-tolerance by using batch processing to provide comprehensive and accurate views of batch data, while simultaneously using real-time stream processing to provide views of online data. The two view outputs may be joined before presentation. The rise of lambda architecture is correlated with the growth of big data, real-time analytics, and the drive to mitigate the latencies of map-reduce.[1]
Lambda architecture depends on a data model with an append-only, immutable data source that serves as a system of record.[2]:32 It is intended for ingesting and processing timestamped events that are appended to existing events rather than overwriting them. State is determined from the natural time-based ordering of the data.

Random thoughts when I play tennis

August 12, 2019

Introduction


It is my personal finance research. I like to write a short research topic called what is my biggest finance mistake made from 2010 to 2019. I did find out in 2019 that I should read my own 401 K statement from Par from 2007 to 2019, and then learn basics of investment, investing on index fund from 2009 to 2019. What is another one?

Another biggest one


I thought about so many times when I played tennis three hours on August 11, 2019. I should learn by myself about technology and system design. I am a single person, and one income family; so it is important for me to reduce the risk of unemployment in Canada; I also try to live under my means, learn how to generate passive income as a landlord to manage a condo in Florida.

I did study the retirement guide and set up a few portfolios to invest the stock market, and grow with USA economy.

I should build a habit to keep learning. One thing I like to learn is to find those technical video related to Facebook, Instagram, Amazon AWS, S3, and I should have learned from 2010 to 2019.

I believe that I should invest my time to learn those technology by myself. Definitely it does not cost me anything. I believe that I can build good habit to learn best engineers in the world, and also help myself to be a good engineer as well.




Streaming Data Solutions on AWS with Amazon Kinesis

Introduction


It is my personal finance research. I try to figure out  what motivates me. How do I push myself to learn in a week to prepare most important onsite interviews related to system design and product design. I am so happy to read the first white paper from Amazon.


Technical detail



27 pages white paper is here to download.

I plan to read the paper in 30 minutes.

You need a different set of tools to collect, prepare, and process real-time streaming data than those tools that you have traditionally used for batch analytics. With traditional analytics, you gather the data, load it periodically into a database, and analyze it hours, days, or weeks later. Analyzing real-time data requires a different approach. Instead of running database queries over stored data, stream processing applications process data continuously in real time, even before it is stored. Streaming data can come in at a blistering pace and data volumes can vary up and down at any time. Stream data processing platforms have to be able to handle the speed and variability of incoming data and process it as it arrives, often millions to hundreds of millions of events per hour.

Case study


ABC tolling company

First requirement:

ABC Tolls would like to make some modifications to its system. The first requirement comes from its business analyst team. They have asked for the ability to run reports from their data warehouse with data that is no older than 30 minutes.

Second requirement:

ABC Tolls is also developing a new mobile application for its customers. While developing the application, they decided to create some new features. One feature gives customers the ability to set a spending threshold for their account. If a customer’s cumulative toll bill surpasses this threshold, ABC Tolls wants to send an in-application message to the customer to notify them that the threshold has been breached within 10 minutes of the breach occurring.


To support the feature to send a notification when a spending threshold is breached, the ABC Tolls development team has created a mobile application and an Amazon DynamoDB table.9 The application allows customers to set their threshold, and the table stores this value for each customer. The table is also used to store the cumulative amount spent by each customer, each month. To provide timely notifications, ABC Tolls needs to update the cumulative value in this table in a timely manner, and compare that value with the threshold to determine if a notification should be sent to the customer. Since their toll transactions are already streaming through Kinesis Firehose, they decided to use this streaming data as the source for their aggregation and alerting. And because Kinesis Analytics enabled them to use SQL to aggregate the streaming data, it is an ideal solution to the problem. In this solution, Kinesis Analytics totals the value of the transactions for each customer over a 10-minute time period (window). At the end of the window, it sends the totals to a Kinesis stream. This stream is the event source for an AWS Lambda function. The Lambda function queries the DynamoDB table to retrieve the thresholds and current total spent by each customer represented in the output from Kinesis Analytics. For each customer, the Lambda function updates the current total in DynamoDB and also compares the total with the threshold. If the threshold has been exceeded, it uses the AWS SDK to tell Amazon Simple Notification Service (SNS) to send a notification to the customers.

Amazon Kinesis Easily collect, process, and analyze video and data streams in real time

Here is the web page to read.


The Log: What every software engineer should know about real-time data's unifying abstraction

Here is the link.

One of the most useful things I learned in all this was that many of the things we were building had a very simple concept at their heart: the log. Sometimes called write-ahead logs or commit logs or transaction logs, logs have been around almost as long as computers and are at the heart of many distributed data systems and real-time application architectures.

To make this atomic and durable, a database uses a log to write out information about the records they will be modifying, before applying the changes to all the various data structures it maintains. The log is the record of what happened, and each table or index is a projection of this history into some useful data structure or index. Since the log is immediately persisted it is used as the authoritative source in restoring all other persistent structures in the event of a crash.

For a long time, Kafka was a little unique (some would say odd) as an infrastructure product—neither a database nor a log file collection system nor a traditional messaging system. But recently Amazon has offered a service that is very very similar to Kafka called Kinesis. The similarity goes right down to the way partitioning is handled, data is retained, and the fairly odd split in the Kafka API between high- and low-level consumers. I was pretty happy about this. A sign you've created a good infrastructure abstraction is that AWS offers it as a service! Their vision for this seems to be exactly similar to what I am describing: it is the piping that connects all their distributed systems—DynamoDB, RedShift, S3, etc.—as well as the basis for distributed stream processing using EC2.

Relationship to ETL and the Data Warehouse


The key problem for a data-centric organization is coupling the clean integrated data to the data warehouse. A data warehouse is a piece of batch query infrastructure which is well suited to many kinds of reporting and ad hoc analysis, particularly when the queries involve simple counting, aggregation, and filtering. But having a batch system be the only repository of clean complete data means the data is unavailable for systems requiring a real-time feed—real-time processing, search indexing, monitoring systems, etc.

Indexing for full text search in PostgreSQL

Here is the article.


Book chapter: The future of Data systems

August 12, 2019

Introduction


It is my favorite book chapter. I have good time to learn something from the book chapter. I just start to read the book chapter less than one week ago, I tried a few times, each time I had to go back to previous chapters, and then learned some basics first. I like to take some notes for each subtitle in the book chapter.

Data Integration


Unbunding Databases


Aiming for Correctness


Doing the right thing

Storage engine - log-structured storage, B-Trees, and column-oriented storage - Chapter 3
replication - single leader, multi-leader and leaderless approaches - Chapter 5

Durable system record - ?
Conversely, search indexes are generally not very suitable as a durable system of record, and so many applications need to combine two different tools in order to satisfy all of the requirements.


Page 492
At an abstract level, they achieve a similar goal by different means. Distributed transactions
decide on an ordering of writes by using locks for mutual exclusion (see
“Two-Phase Locking (2PL)” on page 257), while CDC and event sourcing use a log
for ordering. Distributed transactions use atomic commit to ensure that changes take
effect exactly once, while log-based systems are often based on deterministic retry
and idempotence.

Abstract level
Distributed transactions
ordering of writes by using locks from mutual exclusion
CDC and event sourcing
use a log for ordering

Statement:
Distributed transactions use atomic commit to ensure that changes take effect exactly once, while log-based systems are often based on deterministic retry and idempotence.


In “Aiming for Correctness” on page 515 we will discuss some approaches for implementing
stronger guarantees on top of asynchronously derived systems, and work
toward a middle ground between distributed transactions and asynchronous logbased
systems.


Sunday, August 11, 2019

Memcache Basics

Here is the link.


Reddit Architecture

Here is the link.

Scaling - web development




Problems with Memcachedb - Web Development


Distributed Systems in One Lesson by Tim Berglund

Here is the link.


Normally simple tasks like running a program or storing and retrieving data become much more complicated when we start to do them on collections of computers, rather than single machines. Distributed systems has become a key architectural concern, and affects everything a program would normally do—giving us enormous power, but at the cost of increased complexity as well. Using a series of examples all set in a coffee shop, we’ll explore topics like distributed storage, computation, timing, messaging, and consensus. You'll leave with a good grasp of each of these problems, and a solid understanding of the ecosystem of open-source tools in the space.

Three characteristic

The computers operate concurrently
The computers fail independently
The computers do not share a global clock

Three topics

Storage
Computation
Messaging

Single-master storage

More read than write

Read replication

Distributed database, replicate two databases - use space buy time, short time to read
Eventually consistent database

Master, lead server

10:22

Sharding

Break data model - cannot join cross shards -

How to take joins away? Denormalize, read is slow. To add Index is not option.

Consistent hashing - technique used

Cassandra database -

Replication -> two copies - distributed system, computer fails - I solved the problem, be up, tolerant of hardware failure

Consistency - three copies of data.

A rule -

Consistency

R+W> N

N - number of replicas

Strongly consistent database

19:40/48:59
CAP Theorem

Simple way to understand - distributed database -
Consistency
Available - read works, write works
Partition tolerance

Shared writing project
Coffee shop closes
Synchronizing over the phone
Battery dies
Status report

Choose unavailable - Give you an answer, but it may be wrong
Give up consistency
Sometimes we can not have three

Distributed computation

One process -

MapReduce

Map

Counting words,

Shuffle

similar words near each other

Reduce - add up the number,

300 computers, replicated all

Hadoop

MapReduce API
MapReduce job management
Distributed Filesystem (HDFS)
Enormous ecosystem

Spark

Scatter/gather paradigm (similar to MapReduce)
More general data model (RDDs, DataSets)
More general programming model(transform/action)
Storage agnostic

Kakfa


Everything is stream
Focuses on real-time analysis, not batch jobs
Streams and streams only
Except streams are also tables (sometimes)
No cluster required

Messaging

Means of loosely coupling subsystems
Messages consumed by subscribers
Created by one or more producers
Organized into topics
Processed by brokers
Usually persistent over the short term

Messaging problems

What if a topic gets too big for one computer?
What if one computer is not reliable


Apache Kafka

Definitions
Message: an immutable array of bytes
Topic: a feed of messages
Producer
Consumer
Broker

Message queue

Topic partitioning

0
1
2

3 brokers - one partition

Kafka (Interesting Version)

So you'v got a message bus now
Yesterday you turned events into rows-in-place
Tomorrow maybe you'll just compute events

without streams






Why would a new developer choose Django?

Here is the link.


Turning Caches into Distributed Systems with mcrouter - Data@Scale

Here is the link.


Scaling Memcache at Facebook

Here is the link.

The speaker's linkedin profile is here.


Infrastructure requirements for Facebook
1. near real-time communication
2. Aggregate content on-the-fly from multiple sources
3. Be able to access and update very popular shared content
4. Scale to process millions of user requests per second

Design requirements

Support a very heavy read load
  over 1 billion reads/ second
 insulate backend services from high read rates

Geographically Distributed

memcached

Basic building block for a distributed key-value store for Facebook
 Trillions of items
Billions of requests/second
Network attached in-memory hash table
  Supports LRU based eviction

Roadmap

1. Single front-end cluster
  .Read heavy workload
  .Wide fanout
  . Handling failures
2. Multiple front-end clusters
. Controlling data replication
. Data consistency
3. Multiple regions
 . Data consistency

Why Separate Cache?
High fanout and multiple rounds of data fetching
Data dependency DAG for a small request

Need more read capacity

Problems with look-aside caching
Stale sets

Thundering Herds

Incast congestion





What is Open Graph

Here is the link.

Year 2010
Every one builds their identity by Like activity.

Build identity through Facebook - Time line

Build time line - open graph - tons of more stuff






TAO: The power of the graph

Here is the link.

I plan to read the article from 10:00 AM - 10:110 AM on August 11, 2019.

Even though memcache has "cache" in its name, it's really a general-purpose networked in-memory data-store with a key-value data model.

I need to get familiar with terminologies used in the article

wordload on its data backend
social graph

Users sees
News Feed stories
comments
likes
share for those stories
photos and check-ins from their friends

Best relational database technology
a poor match
supplemented by a large distributed cache that offloads the persistent store.

Bugs
User-visible inconsistencies
site performance issue

product engineer
two data stores
different data models: a large cluster of MySQL servers for storing data persistently in relational tables
an equally large collection of memcache servers for storing and serving flat key-value pairs derived.

Implementation

The TAO service runs across a collection of server clusters geographically distributed and organized logically as a tree. Separate clusters are used for storing objects and associations persistenly, and for caching them in RAM and FLASH memory. This separation allows us to scale different types of clusters independently and to make efficient use of the server hardware.

Caching clusters running TAO servers - what are caching clusters?
In addition to satisfying most read requests from a write-through cache, TAO servers orchestrate the execution of writes and maintain cache consistency among all TAO clusters. We continue to use MySQL to manage persistent storage for TAO objects and associations.

The data set managed by TAO is partitioned into hundreds of thousands of shards. All objects and associations in the same shard are stored persistently in the same MySQL database, and are cached on the same set of servers in each caching cluster. Individual objects and associations can optionally be assigned to specific shards at creating time. Controlling the degree of data collocation proved to be an important optimization technique for reducing communication overhead and avoiding hot spots.

Shards can be migrated or cloned among servers in the same cluster to equalize the load and to smooth out load spikes.

Actionable Items


I like to read the article and record using my Samsung phone. And then I like to share my recording as an attachment.


Profiling Python with cProfile

Here is 2 minutes read article.


Carl Meyer about Django @ Instagram at Django: Under The Hood 2016

Here is the link.

Django: Under The Hood is an annual Django conference for experienced Django developers. Come and learn about the internals of Django, and help to shape its future.

AppWeight

IG Weight
Instagram

33:00/ 1:04:32
Scalability - measure

Metrics - visible to engineers

Dynostats

CProfile

Fixing efficiency regressions

42:00 - very good talk, give examples of problems found:

Fix the obvious
Don't do useless work
Cache things that don't change - do it once
Change ...



Questions:
50:00
1. Profile tools?
2. Monitoring tools?
3. TAO vs moderate website?
4.


The big interview with Martin Kleppmann: “Figuring out the future of distributed data systems”

Here is the link.


What is a Container?

Here is the link.




Containers vs. Virtual Machines (VMs): What’s the Difference?

Here is the link.


Saturday, August 10, 2019

Containers and VMs - A Practical Comparison

Here is the link.


Project Hatchway: Persistent Storage for Cloud-Native Applications

Here is the article to read.


Persistent Container Storage with Project Hatchway

Here is the link.


Kubernetes

Here is the wiki page for Kubernets.


Kubernetes in 5 mins

Here is the five minutes video presented by Steve Tegeler. The author's linkedin profile is here.

Kubernetes in 5 min

"Desired State managment"

K8S cluster service -> Worker

App1.yam1

Deployment
Pod1
-contimg1
-contimg2

Replica -> 3
Pod2
  container3
  replicas = 2


Jay Kreps | Kafka Summit SF 2018 Keynote (Kafka and Event-Oriented Architecture)

Here is the link.

Overcome by Events (OBE)

A term of military origin used
What is an event? Something to occur

Where are they?
Events haven't had a proper home in infrastructure of in code.

Microservices
Monitoring
Data pipelines
Analytics

7:32/ 24:34

Monitoring


Old School
4 Pillars
. Metrics
. Alerting
. Distributed Tracing
. Log Aggregation

New School
Event processing

Log aggregation

Data pipelines



Analytics 



The stream platform


01 Immutable log-centric architecture
02 Support for tables with log compaction
03 Horizontal scalable
04 Connector architecture
05 Transactional semantics for stream processing

01 Immutable log-centric architecture
abstraction - powerful idea -

02 Support for tables with log compaction

03 Horizontal scalable



Ask Confluent #6: Kafka, Partitions, and Exactly Once ft. Jason Gustafson

Here is the link.


Martin Kleppmann

Here is the link.


Top 10 engineers I choose to be my mentor this August

August 10, 2019

Introduction


I am super powerful programmer, but I am not a good thinker in terms of distributed system design or product designer, what can I do? I like to find top 10 engineers and then I like to learn from their teaching.

Top 10 engineers



I like to put down the list of top 10 engineers. Since I only have 10 days to prepare since I take vacation. I have to count down hours, and I like to challenge myself how smart I can, learn from those top engineers.

1. Guo Lisa, scale instagram infrastructure, here is the video.
2. Martin Kleppmann  28 minutes video is here.
3. Tim Berglund, Distributed Systems in One lesson by Tim Berglund, the link is here.
4. Udacity course - reddit architecture (Added on 8/13/2019 11:22 PM)
Web development, the link is here.
My favorite topics are Reddit architecture, Reddit architecture 2.


Martin Kleppmann | Kafka Summit SF 2018 Keynote (Is Kafka a Database?)

Here is the link.

Kafka Summit SF 2018 Keynote by Martin Kleppmann (Researcher, University of Cambridge). Martin Kleppmann is a distributed systems researcher at the University of Cambridge, and author of the acclaimed O’Reilly book “Designing Data-Intensive Applications” (http://dataintensive.net/). Previously he was a software engineer and entrepreneur, co-founding and selling two startups, and working on large-scale data infrastructure at LinkedIn.


ACID - Atomic, Consistency, Integrity, Durability

Durability

3:05 Kafka - provide durability
Write to disk, async
Archive tape, fsync to disk
Copy to several machines, if you lose the machine, then you will still have the data.
Atomicity - Handling faults (crashes)

Concurrency

Write data to different places, do not let web app directly write to those places. Write to Kafka.

The log of Karfa, guarantee now ..., atomic is easy to write a single event, a stream.

They independently consume the log.

6:55 PM
13:57/ 28:14

How to do it in Karfka?

Put an event in Karfka, Json to document the transaction.

{eventType: transfer, from Account: 12345, toAccount: 54321, amount: 100.0, event ID:abcdef}

Two events: take 100 from the account, input 100 into the account

Second stream process -



Isolation - ACID

Serializable - as if there is a database available, ...

relational database:

start transaction

select count(*)
from user_account
where user = 'jane'

Non-serializable execution

Consistency - ACID

Enforcing invariants
Integrity ...

Sooo... is this a database?
No ad-hoc queries, for sure ...

26:50
Transactions broken down into multi-stage stream pipelines

But: stronger consistency properties than many distributed datastores...




How should I choose who to learn from?

August 10, 2019

Introduction


It is my personal finance research. I like to take some time to learn from best in the industry to prepare my two onsite interviews in next two weeks, Amazon and Facebook. I like to figure out who I like to choose to be my coach to learn some basics from their sharing and presentation in the conference.

How should I choose who to learn from?


It is very challenging task for me to expedite my learning process. I need time to warm up algorithm and data structure problem solving, but I also like to push myself to learn how to be a designer.

I am a single person. I feel lonely all the time, so I will try to reach out wechat to check messages or take a break to play tennis sports. But I think that there are so many good videos from youtube.com, I should spend time to catch up learning for my onsite interviews.


Get to know the industry

August 10, 2019

Introduction


I have system design and product design interview. I have a book to read, and also videos to watch related to system design. But also I like to spend time to learn some those talent engineers and watch a few videos to get educated quickly.

First 10 videos count down


The first video is from Facebook Guo lisa, here is the blog.

I like to have 10 videos to watch, and then I like to count down the number.

I wish that I have watched thousands of videos related to technologies from 2010 to 2019 August. But in reality, I was too busy to do so many things. I never learn that it is so important for me to reach out, and set the role models from those engineers, leaders who share the technology video on youtube.com




Lessons learned form Kafka in production (Tim Berglund, Confluent)

Here is the link.

Many developers have already wrapped their minds around the basic architecture and APIs of Kafka as a message queue and a streaming platform. But can they keep it running in production? This talk contains real-world troubleshooting and optimization scenarios culled from the logs of Confluent technical support. We’ll talk about the trade-offs between optimizing for the always-desirable outcomes of throughput, latency, durability, and availability. How many partitions should you use for a given topic? How much message batching should you configure in the producer? How many replicas should be required to acknowledge a write? What do you do when you see a partition growing inexplicably? When should you rearchitect your application to use the streaming API? We’ll answer these questions and more int his overview of common Kafka production issues.

Outline
1. A reminder about streams
2. Apache kafka review
3. Strange happenings with partitions
4. Automated liveness check
5. Adding a broker hurts!

Streaming platform

Event-centric thinking

Mobile    ---------->
Web app   --------->    streaming platform  ->  Rec engine, security, hadoop, real-time analytics
API          ---------->

Forward compatible

Outline

Producers -> Karfka cluster -> consumers

Logs
0 1 2 3 4 5 6 7 8 9  <-------------  next write
            |___  reader

KAFKA Topic = Partitioned Log

Order is only in each partition. But ......

API level

30:44

Partition?
. String consistency
. Leader
. Follower(s)
. Controller


three replications, 4 partitions


34:56
In-sync replicas
.Replica can fall behind
.ISR list grows, shrinks
.What if it gets too small?

Moral:
Waht your ISR list!

41:23
. > partitions, > throughput
. Plan ahead
. Clean leader shutdown
. Leader failure

Moral
Automation is a sharp knife

Auto data balancing


Actionable Items


I really enjoy the learning from the presentation. I like to watch a few more videos so that I can learn the rough idea how to manage my preparation.




GOTO 2016 • What I Wish I Had Known Before Scaling Uber to 1000 Services • Matt Ranney

Here is the link.

To Keep up with Uber's growth, we've embraced microservices in a big way. This has led to an explosion of new services, crossing over 1,000 production services in early March 2016. Along the way we've learned a lot, and if we had to do it all over again [...]

5:31
Microservice
immutable?
Append only?

Why microservices?

Now you have a distributed system
Eveything is an RPC
what if it breaks?

Less obvious costs
Everything is a tradeoff
You can build around problems
Might trade complexity for politics
You get to keep your biases

Languages

Hard to share code
Hard to move between teams
WIWK: Fragments the culture

Fragment the culture - human behavior - really cost

RPC
HTTP/REST gets complicated
JSON needs a schema
RPCs are slower than PCs
WIWIK: servers are not browsers


Performance
Doesn't matter until it does
Probably want at least simple perf requirements
WIWIK: "good" not required, but "known" is

FANOUT

overall latency >= latency of slowest
1ms avg, 1000ms p99
use 1:1% at least 1000ms
use 100: 63% at least

Tracing
Lots of ways to get this
Best way to understand fanout

Tracing 
Probably want sampleing
WIWIK: cross-lang context propagation

LOAD TESTING

Need to test against production
Without breaking metrics
Preferably all the time
WIWIK: all systems need to handle "test" traffic

MIGRATIONS
old stuff still has to work
What happened to immutable?
WIWIK: mandates are bad







Friday, August 9, 2019

CoDel schedule algorithm

Here is the link.


Keynote - Systems at Facebook Scale

Here is the link.

Tech-lead of web foundation team

Facebook in 30 seconds



Rapid change

- Code released twice a day
- Rapid feature development - e.g. Lookback videos
 - 450 Gbps of egress
 - 720 million videos rendered ( 9 million/ hour)


Proactive design for scaling

Solve scaling problems only once
All Facebook projects in a single source control repo - easy code reuse

Folly: base C++ library
Thrift: RPC
Proxygen: HTT(s) server
RocksDB: persistent key-value store

Efficient synchronization

Producer/Consumer queues

How to implement Producer/Consumer

Multiple wakeups/ Deque

Potential context switches

Word thread ordering
FIFO
LIFO

Good queueing
Deal with a burst of load
Processing speed > arrival speed
Increase reliability

Bad queueing
Server is overloaded
Processing speed < arrival speed
Causes latency

Overload handling philosophy
If you're going to fail, fail quickly - prevents delaying the overall request
Servers should signal overload - OK to say "I can't help you right now"
Order doesn't matter - servers need not first-com-first-serve
Clients should defend themselves - don't rely on the server
complex knobs are tuned poorly - design parameter-free abstractions

Controlled delay






Advanced Message Queuing Protocol

Here is wiki article.