Sunday, March 27, 2022

Ming Dao school: Mock interviews

 Ming Dao School conducts system design mock interviews online every Sunday. You are welcome to join us.

Here is the link. 

March 27, 2022

Mock interview - Senior engineer -  system design 

Requirements


Functional requirements

Inventory management like instacart

Make sure we don’t overshop

Shipment coming in

Browse, add to cart, at the same time update availability, no over selling


Out of scope

No payment, shipment


  1. Inventory service, we already have item in

  2. Neighboring team - catalog team. Already handle catalog

  3. Manage inventory and basket


0. Incoming shipment to add to inventory

1. User can browse, search. Search and browse for bare minimum

2. Add to cart, keep something (lock/reserve)

(Really want to focus on inventory management)


Non functional requirements

Availability

High throughput

Service reliability


System Design


10:00


External APIs


API:

  1. add_product(product, number)

  2. Add_to_cart(user, product, number)

  3. checkout(user, product, number)


11:00

Database design


User_purchse_table

User_id, product_id, number, status

Status: in_cart, bought


Product_table

Product_id, sku_number, total_number


Single service logic:


Add_to_cart: do db txn with 2 modification, and commit


For transaction: mysql. May have problem, the size may be too big


15:00

A truck coming in with a box of banana.  There is a person with scanner.


What does “add product” do?  Yes. we can call add product for the above situation


If we have never sold iPad, now we start to sell iPad.  


Maybe register product. But not sure distinction


User_purchase_table

product_table


Sharding with MySQL can be hard

Instead of one transaction to deal with 2 tables, we can do something different.


Server logic:

Modify product_table if products number >= N

If server logic succeed:

Modify user_purchase_table, then return success

Else:

Return true


Q: When you shard it. If it’s same db is fine to use transaction support.

Now we are using distributed transaction.


A: Break transaction into 2 phases.


Various crash scenarios, because data relevant to the transaction are separate in 2 databases


We have idempotent key for user behavior.  If first step succeed, second step may fail


Idempotent key has table: idempotent_key, status


If we need to give up after we reserve banana.  How do we rollback?


  1. Have global key for idempotent key

  2. Chron job to periodically check: consistency between idempotent_table and product_table

Treat as idemponent key: unique id: lock N products when number of products equal>=N. CAS

Are all APIs on the same server? Who is calling those APIs?


If there are many stores, and they are busy.  Talk me through sharding


Most important is the product_table


How do you shard?

We can shard based on category and product name



Downside: read-write heavy. Not good for hot products


We can cache the product



Interviewer and Audience Feedback


Audience feedback


Inventory management for grocery

Instacart: cannot sell more than the inventory.

Consistent.



Senior:

Leading the discussion

Which points are most important

Important points:

  • Managing inventory

  • Abandoned cart, expiry.  A hidden requirement


Transaction - consistency

Possible:

  • Distributed transaction

  • 2PC, saga

  • Idempotent

  • Can proactively provide the possible solutions

  • Roll-forward, rollback.

  • Touched on the points

  • Most of these are based on data.  Today, we are missing data estimation

  • Everybody has reservation.  How many update.  What’s the granularity update?  Can we use in-memory database? 

  • Which one is better?

  • Relational database.  Every store, every item, may not have a lot of entries.

  • We can shard on store.

  • Using business requirement to optimize


Shopper fulfillment and inventory

  • Shopper: more API. 

  • Should use picture


Shipment acceptance:

  • spike 

  • There are multiple solutions


Some may be very specific. But big companies just look at big picture but not specific skills


Requirement gathering.  We can simplify the solution, senior should drive simplification

Finding the key points of discussion.  QPS. storage. Hit rate, locality. Prove your own design.  

Present the big picture.  Search can be separated from this system. 


Weak pass for intermediate level


Interviewee


Q: Difference between adding a new item, vs adding a new product.

A: Add item to category is separate from add item to inventory.


Q: sharding. We can use SQL. For distributed transaction, how do we do?

A: NoSQL.

We can do distributed transaction. E.g. 

Orchestrator for distributed transaction

Complexity in team.

2-phase commit is relatively slow

QPS is not going to be very high


SAGA:  transaction log, decide where we are. Roll-forward or roll-backward.

Need to write transaction log. 

Inventory table - need versioning.


Need orchestrator, or transaction log.


Justify relation.  But we should do some calculation

Otherise, we can just subtract number from everyone’s cart

扬长避短


Distributed transaction, 2PC, SAGA


There may be service frontend for two services.


Product requirement

Why do we need to lock

Sometimes we may use substitute

Black store: we have more control.  Substitution - product impact

This is more related to black store / warehouse


Add cart, reserve.


Interviewee:

A: How to do cache?

B: traffic pattern, read/write ratio, locality, hot/cold data

Justification. 

Write-aside, write through

Which system writes to cache.

Invalidate.  Can provide more details


Interviewee:

Biggest impression: distributed transaction is hard.  We can take out some framework

We can provide some transaction methods

E.g. dynamoDB can provide distributed transactions


Most popular:

Saga


Assume distributed transaction can solve this.

If we don’t use RDBMS, what can we discuss?

We can discuss expiry.  In-memory.  We don’t have to compute the total number.


You can lead the discussion as a senior engineer.


2PC, spanner.

Before diving to details,draw the big picture

Need to draw UI


Should draw the UI

Bookkeeping are important


Dive too deep into transactions


Interviewer: thinks mysql can solve the problem

Sharding, replication, index.


Q: distributed transaction

2PC vs saga

It’s a design pattern.  Roll forward, rollback.

It does not guarantee isolation.  Temporarily the data is not right.


When we buy things, do we save this to backend cache?

Load balance: backend server.  How do we memorize these?


We ideally want to set up stateless server

Redis, elastic cache.  

Redis: single instance is the strongest.  If you need durability, you can use replicated/distributed redis.


Single master, single instance, quorum, paxos, raft.  Leader election mechanism.


Is it cache or storage? Redis is quite durable.


If we use RDBMS, then things are simple.

If there is no locality and miss rate is high, then cache is not useful.


Sticky routing: fairly rare.


Complexity, hardware cost, team cost. 

We can use long polling.  We don’t have to use realtime feedback.  Over-engineering


Cache or database.

High speed cache.

All carts have expiration.

In-memory database + replication (providing durability)


Every shop does not have a lot of people, so we can calculate in realtime.


Cache -> roadmap with transaction


=

Put all order in cart

userID + cart

Cart may have item A

Cart -> split into user, item, store, cart, reservation number. Multiple entries.

Reservation count storage

We don’t need to worry about expiration.


TTL.  How to add the reserved number?

If we see TTL expired, we can remove the item, and not add the count


Alternatively, we can use a sleeper to poll. Add the things back.



Reservation table, primary key?

Store, SKU, cart ID, product ID, reservation count

Primary key is concatenation of the 4

Cart: should add based on store, SKU.  Range scan.

Black warehouse.  Popular items.  Say apple is most popular.


Can we set store, sku as primary key, the rest are repeated.

It makes things more complex: Normalization, denomalization. Locking, optimistic locking.


DynamoDB is expensive: you may as well use RDBMS. RDBMS can be tuned as well.


==

If we store and shard by store, then retrieval is hard

You need to go through the shards.

We can materialize view.  Enrichment.  Precalculate the additional information.


Can we use cart as primary key?

Cart + map

Cart + reservation number


Secondary index: equivalent of more tables


Global secondary index: internally there are 2 tables. It’s eventual consistency.  It uses distributed transaction.


如何准备?

DDIA - backburner

Alex Xu - 2 books

Tech Dummies - it gives a brief introduction

Come to mock

以面代考


Pattern of distributed system

DDIA

New book - stock exchange, payment, google map

If we hit a product that you have not used.


If you don’t know the design at all:

Start from user, end with user

Big picture design

Find the things that I am familiar with, and discuss about it. E.g. how to split the microservice


15 companies

Mock, listen, read, real interview


Grokking: relatively shallower

Alex Xu: deeper, but still some area not deep enough

Mock with each other: can get real architecture


==


Excalidraw, googledocs

Multiple cameras

Open tethering

Excel

Work from alpha

Open your IDE


Debugging, coding

Coding more 

After SDE II, SDE III 

Example: organization level. Executive. Data point

Cross functional. L6


Cross functional


Situation of conflict, example requirements of the person from the other side, how we pull the other side back.




Friday, March 25, 2022

Apache Solr

 Solr is an open-source enterprise-search platform, written in Java. Its major features include full-text search, hit highlighting, faceted search, real-time indexing, dynamic clustering, database integration, NoSQL features and rich document handling.

Udemy: System design | My purchase


My new purchase



Thursday, March 24, 2022

Design Facebook Messenger

 Here is the link.

Watch Jacob, Exponent co-founder and Dropbox software engineer, answer the system design question "Design Facebook Messenger"

To submit your own answer for this question and get feedback, visit the answer in our interview question database.

Kubernetes | 10 minutes study

 A Kubernetes cluster consists of a set of worker machines, called nodes, that run containerized applications. Every cluster has at least one worker node.

The worker node(s) host the Pods that are the components of the application workload. The control plane manages the worker nodes and the Pods in the cluster. In production environments, the control plane usually runs across multiple computers and a cluster usually runs multiple nodes, providing fault-tolerance and high availability.

This document outlines the various components you need to have for a complete and working Kubernetes cluster.


Terraform (software)

 Terraform is an open-source infrastructure as code software tool created by HashiCorp. Users define and provide data center infrastructure using a declarative configuration language known as HashiCorp Configuration Language (HCL), or optionally JSON.[3]

Terraform manages external resources (such as public cloud infrastructure, private cloud infrastructure, network appliances, software as a service, and platform as a service) with "providers". HashiCorp maintains an extensive list of official providers, and can also integrate with community-developed providers.[4] Users can interact with Terraform providers by declaring resources[5] or by calling data sources.[6] Rather than using imperative commands to provision resources, Terraform uses declarative configuration to describe the desired final state. Once a user invokes Terraform on a given resource, Terraform will perform CRUD actions on the user's behalf to accomplish the desired state.[7] The infrastructure as code can be written as modules, promoting reusability and maintainability.[8]

Terraform supports a number of cloud infrastructure providers such as Amazon Web Services, Microsoft Azure, IBM Cloud, Serverspace, Google Cloud Platform,[9] DigitalOcean,[10] Oracle Cloud Infrastructure, Yandex.Cloud,[11] VMware vSphere, and OpenStack.[12][13][14][15][16]

HashiCorp also supports a Terraform Module Registry, launched in 2017.[17] In 2019, Terraform introduced the paid version called Terraform Enterprise for larger organizations.[18]

Terraform has four major commands:

$ terraform init
$ terraform plan
$ terraform apply
$ terraform destroy

全球公有云编排服务大比拼

 【摘要】 本文深入并客观地分析全球主要的公有云平台中,资源编排服务/应用编排服务的能力象限图。包括AWS的Cloudformation,阿里的ROS,OpenStack的Heat,华为的AOS,微软的RM,谷歌的CDM,以及青云的RO,K8S生态中的Helm,甚至对腾讯云的编排能力也做了调侃。

领域澄清

首先明确我们在讲什么,也就是本文描述的应用编排,资源编排是什么东西。在云上编排的含义一般有两种:

1. 云平台上自动化创建云服务,并部署应用。叫做资源编排or应用编排or服务编排。

2. 容器应用,根据资源要求,调度到哪个节点上。叫做容器编排or资源调度。

Tricks of the Trade: Tuning JVM Memory for Large-scale Services

 

Running queries on Uber’s data platform lets us make data-driven decisions at every level, from forecasting rider demand during high traffic events to identifying and addressing bottlenecks in the driver sign-up process. Our Apache Hadoop-based data platform ingests hundreds of petabytes of analytical data with minimum latency and stores it in a data lake built on top of the Hadoop Distributed File System (HDFS). 

Our data platform leverages several open source projects (Apache Hive, Apache Presto, and Apache Spark) for both interactive and long running queries, serving the myriad needs of different teams at Uber. All of these services were built in Java or Scala and run on open source Java Virtual Machine (JVM).  

Uber’s growth over the last few years exponentially increased both the volume of data and the associated access loads required to process it, resulting in much more memory consumption from services. Increased memory consumption exposed a variety of issues, including long garbage collection (GC) pauses, memory corruption, out-of-memory (OOM) exceptions, and memory leaks. 

Refining this core area of our data platform ensures that decision-makers within Uber get actionable business intelligence in a timely manner, letting us deliver the best possible services for our users, whether it’s connecting riders with drivers, restaurants with delivery people, or freight shippers with carriers. 

Preserving the reliability and performance of our internal data services required tuning the GC parameters and memory sizes and reducing the rate at which the system generated Java objects. Along the road, it helped us develop best practices around tuning the JVM for our scale which we hope others in the community will find useful.

What is JVM garbage collection?

The JVM runs on a local machine and functions as an operating system to the Java programs written to execute in it. It translates the instructions from its running programs into instructions and commands that run on the local operating system. 

The JVM garbage collection process looks at heap memory, identifies which objects are in use and which are not, and deletes the unused objects to reclaim memory that can be leveraged for other purposes. The JVM heap consists of smaller parts or generations: Young Generation, Old Generation, and Permanent Generation. 

The Young Generation is where all new objects are allocated and aged, meaning their time in existence is monitored. When the Young Generation fills up, using its entire allocated memory, a minor garbage collection occurs. All minor garbage collections are “Stop the World” events, meaning that the JVM stops all application threads until the operation completes. 

The Old Generation is used to store long-surviving objects. Typically, a threshold is set for each Young Generation object, and when that age is met, the object gets moved to the Old Generation. Eventually, the Old Generation needs a major garbage collection, which can be either a full or partial “Stop the World” event depending on the type of garbage collection configured in the JVM program arguments. 

The Permanent Generation stores classes or interned character strings. It is not for objects that survived from the Old Generation to stay permanently. If this area is about to be full, there will be a GC, which is still counted as a major GC. 

The JVM garbage collectors include traditional ones like Serial GC, Parallel GC, Concurrent Mark Sweep (CMS) GC, Garbage First Garbage Collector (G1 GC), and several new ones like Zing/C4, Shenandoah, and ZGC. 

The process of collecting garbage typically includes marking, sweeping, and compacting phases, but there can be exceptions for different collectors. Serial GC is a rudimentary garbage collector, which stops the application for the whole collecting process and collects garbage in a serial manner. Parallel GC does all the steps in a multi-threaded manner, increasing the collecting throughput. CMS GC attempts to minimize the pauses by doing most of the garbage collection work concurrently with the application threads. The G1 GC collector is a parallel, concurrent, and incrementally compacting low-pause garbage collector. 

With these traditional garbage collectors, the GC pause time usually increases when the JVM heap size is increased. This problem is more severe in large-scale services because they usually need to have a large heap, e.g., several hundreds of gigabytes. 

The new garbage collectors like Zing/C4, Shenandoah, and ZGC try to solve this problem, minimizing the pauses by running the collecting phases concurrently and incrementally more frequently than traditional collectors. Also, the pause time doesn’t increase as the heap size goes up, which is desirable for large scale services within a data infrastructure.   

In addition to garbage collectors, the object creation rate also impacts the frequency and duration of GC pauses. Here, the object creation rate defines how many objects with size in bytes are created for a given time range in seconds. It is easy to understand that when the rate goes higher, more objects are created and occupy the heap, triggering the GC more frequently and causing longer pauses.

In our experience with garbage collection on JVM, we’ve identified five key takeaways for other practitioners to leverage when working with such systems at Uber-scale.

Tricks of the Trade: Tuning JVM Memory for Large-scale Services

March 24, 2022

Here is the article. 

Key takeaways

Through our experience maintaining and improving query support for Uber’s data platform, we learned many critical lessons about what it takes to optimize successful JVM memory and GC tuning. We consolidated these learnings into more granular takeaways for improving JVM and GC performance at Uber’s scale, below:   

  1. Discern if JVM memory tuning is needed. JVM memory tuning is an effective way to improve performance, throughput, and reliability for large scale services like HDFS NameNode, Hive Server2, and Presto coordinator. When GC pauses exceeds 100 milliseconds frequently, performance suffers and GC tuning is usually needed. In our GC tuning scenario, we saw HDFS throughput increase ~50 percent, the HDFS latency decrease ~20 percent, and Presto’s weekly error rate drop from 2.5 percent to 0.73 percent.
  2. Choose the right total heap size. The total JVM heap size determines how often and how long the JVM spends collecting garbage. The actual size is shown as the maximum memory footprint in the verbose GC log. Incorrect heap size could cause poor GC performance and even trigger an out-of-memory exception. 
  3. Choose the right Young Generation heap size. The Young Generation size should be determined after the total heap size is chosen. After setting the heap size, we recommend benchmarking the performance metric against different Young Generation sizes to find the best setting. Typically, the Young Generation should be 20 to 50 percent of the total heap size, but if a service has a high object creation rate, the Young Generation’s size should be increased. 
  4. Determine the most impactful GC parameters. There are many GC parameters to tune, but usually changing one or two parameters makes significant impact while others are negligible. For example, we found changing Young Generation size and ParGCCardsPerStrideChunk improved the performance significantly, but we did not see much difference when changing TLABSize and ConcGCThreads. In tuning Presto GC, we found the String Deduplication setting dominates the performance impact. 
  5. Test next generation GC algorithms. In a large scale data infrastructure, critical services usually have very large JVM heap sizes, ranging from several hundreds of gigabytes to terabytes. Traditional GC algorithms have trouble handling this scale and experience long GC pause times. Next generation GC algorithms, such as C4, ZGC, and Shenandoah, show promising results. In our case, we saw that C4 reduces latency (P90 ~17ms) compared with CMS (P90 ~24ms). 

Moving forward

While we made great progress improving our services for performance, throughput, and reliability by tuning JVM garbage collection for a variety of large-scale services in our data infrastructure over the last two years, there is always more work to be done.  

For instance, we began integrating C4 GC into our HDFS NameNode service in production. With the encouraging and beneficial performance improvement in the staging environment as described above, we believe C4 will help prevent NameNode bottleneck issues and reduce request latency.  

GC tuning in distributed applications, particularly in Apache Spark, is another area we want to examine in the future. For example, the ingestion pipelines used in our data platform are built on top of Spark, and our Hive service also relies on Spark. JVM Profiler, an open source tool developed at Uber, can help us analyze GC performance in Spark so we can improve its performance. 

Forbes: Exponent

 Since 2018, Stephen Cognetta and Jacob Simon (pictured left) have bootstrapped a profitable, San Francisco-based business, Exponent, which offers subscription-based educational and interview-prep tools for aspiring software engineers and product managers. Institutions including Stanford, Yale and Duke have partnered with Exponent to provide career guidance to their students. The startup's YouTube channel has attracted 130,000 followers and 5 million video views and its live practice tools host 12,000 interviews per month.

Forbes Lists

Exponent: Day 2 to be a student | Fundamentals of system design

 March 24, 2022

Synchronous vs asynchronous processes

  1. Batch processing 
  2. Stream processing
  3. Lambda architecture
  4. Asynchronous queues

Batch processing - most popular implementation - MapReduce

Stream processing - data flows into your system as events occur. Events may be changes in state, or updates. 

Asynchronous queues

  1. Message queues
  2. Task queues
  3. Publish/Subscribe (or pub/sub)
Publish/Subscribe (or pub/sub) messaging is another async communication method to know. In pub/sub messaging, you have a subscriber who receives a message sent by a publisher via a broker. Because communication is decoupled, messages can be automatically pushed to all subscribers rather than pulled individually via a message queue. This method is popular in event-driven architecture because event-driven service can be delivered quickly and easily. Additionally, because publishers are isolated from subscribers, the system is easier to maintain and secure. 

When to bring it up in an interview

Consider async processing for applications with:

  • Long expected processing times, or when the processing time is unclear.
  • No need for immediate processing. A classic example is Facebook's newsfeed. A newly uploaded post will be immediately visible in your own feed but it may take some time before it's visible to your entire network.

Batch processing is useful when you need to process chunks of data in predictable intervals, as in accounting software. Stream processing is a good choice for applications like anomaly detection or sentiment analysis when timeliness is critical, and lambda architecture can help you capture the benefits of both if you're willing to deal with added complexity. Synchronous processing is still the best choice for simple applications where processing times are short and well-defined, and when errors must be dealt with immediately like in payment tools.

Take a look at the below table to help guide your decisions when prepping for your interview:

Applications

  • Analytics
  • Web crawling
  • Handling large file uploads
  • Handling real-time events
  • Generating a newsfeed
  • Scheduled tasks
Processing Methods to consider
  • Analytics - Batch processing, MapReduce
  • Web crawling - Batch processing, MapReduce
  • Handling large file uploads - Job queue
  • Handling real-time events - Stream processing
  • Generating a newsfeed - Job queue, Pub/Sub
  • Scheduled tasks - Job queue, Batch processing


MapReduce: Fundamentals of MapReduce with MapReduce Example

March 24, 2022

Here is the example. 

So, just like in the traditional way, I will split the data into smaller parts or blocks and store them in different machines. Then, I will find the highest temperature in each part stored in the corresponding machine. At last, I will combine the results received from each of the machines to have the final output. Let us look at the challenges associated with this traditional approach:

  1. Critical path problem: It is the amount of time taken to finish the job without delaying the next milestone or actual completion date. So, if, any of the machines delay the job, the whole work gets delayed.
  2. Reliability problem: What if, any of the machines which are working with a part of data fails? The management of this failover becomes a challenge.
  3. Equal split issue: How will I divide the data into smaller chunks so that each machine gets even part of data to work with. In other words, how to equally divide the data such that no individual machine is overloaded or underutilized.
  4. Single split may fail: If any of the machines fail to provide the output, I will not be able to calculate the result. So, there should be a mechanism to ensure this fault tolerance capability of the system.
  5. Aggregation of the result: There should be a mechanism to aggregate the result generated by each of the machines to produce the final output.

These are the issues which I will have to take care individually while performing parallel processing of huge datasets when using traditional approaches.

To overcome these issues, we have the MapReduce framework which allows us to perform such parallel computations without bothering about the issues like reliability, fault tolerance etc. Therefore, MapReduce gives you the flexibility to write code logic without caring about the design issues of the system.