Tuesday, August 11, 2020

Netflix tech talk: You Won't Believe How the Biggest Sites Build Scalable and Resilient Systems!

 Here is the link. 

First 28 minutes, sharding, consistent hashing, hotspot, going multi-zone, Cassandra. 

19:50

Sharding

- split writes across master databases

- Each can have a salve, some many slaves based on workload

- One can avoid reading from the master if possible

- Picking the sharing key well is essential and fraught with peril

Second class users

- logged out users get cached content

- CDN bears the brunt of the traffic

Building a data model

- what questions you want to ask your dta?

- Don't try and normalize anything

- Instead of changing a value keep a record of what happened

Data schemas 

- Unless you are really really sure of your business model...

- The less schema the better

- reddit's database is literally just keys and values, despite being in Postgress

28:00 - viewing data, second person to present architecture - Experience evolution, 14 minutes videos. 

I like to watch those 14 minutes 10 times. I need to get ideas how to talk about system design related to this streaming data service, how to scale etc. 

Why Cassandra?

- Availability over consistency

- Writes over reads

- We know Java

- Open source + support


Subscribers - 

Virtuous cycle - viewing -> improved personalization - 

Viewing data 

who, what, when, where, how long 

Real time data use cases

What have I watched?

Where was I?
What else am I watching?

Session Analytics - 

buffering, quality - Session analytics 

Generic architecture 

Architecture evolution

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

Real time data - gen 2 mintivation

Scalability - scale out not up 

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 - ...

gen 3 - requests scale 

Real time data - redistributed...

Stateless microservcies



System design: How Netflix Serves Viewing Data?

 Here is the link. 

Here, viewing data means…

  1. viewing history. What titles have I watched?
  2. viewing progress. Where did I leave off in a given title?
  3. on-going viewers. What else is being watched on my account right now?

The viewing service has two tiers:

  1. stateful tier = active views stored in memory

    • Why? to support the highest volume read/write
    • How to scale out?
      • partitioned into N stateful nodes by account_id mod N
        • One problem is that load is not evenly distributed and hence the system is subject to hot spots
      • CP over AP in CAP theorem, and there is no replica of active states.
        • One failed node will impact 1/nth of the members. So they use stale data to degrade gracefully.
  2. stateless tier = data persistence = Cassandra + Memcached

    • Use Cassandra for very high volume, low latency writes.
      • Data is evenly distributed. No hot spots because of consistent hashing with virtual nodes to partition the data.
    • Use Memcached for very high volume, low latency reads.
      • How to update the cache?
        • after writing to Cassandra, write the updated data back to Memcached
        • eventually consistent to handling multiple writers with a short cache entry TTL and a periodic cache refresh.
      • in the future, prefer Redis’ appending operation to a time-ordered list over “read-modify-writes” in Memcached.

Original Post

I like to watch the video - 50 minutes - 

You Won't Believe How the Biggest Sites Build Scalable and Resilient Systems!

System design: Introduction to Neo4j and Graph Databases

 Here is the link. 

I like to spend next two hours to study the talk. Now it is 1:32 PM, Tuesday, August 11, 2020. 

I also like to take some notes as well. 


Dr. David Allen - 


Go over the example graph from 1:37 - 1:56, and I like to take a dive on example, and see how good it is to learn from Dr. David Allen. 

1:43 

Question and answers:

Developer pages

new4j.com/developer

Intro to Graphs

RDBMS to Graphs

Data Import & Data Modeling

Language Guides & Drivers

Visualization

Neo4j Ecosystem 




Google search: Graph database

 August 11, 2020

Introduction

I like to take some time to look into graph database as my research topic. I may be able to complete a few things in next few hours. 

Graph database

I will work on a few things based on Google search Graph database. 



System design: TAO - graph data store

 August 11, 2020

Introduction

It is my learning process about TAO. I like to take some notes from PPT file of the tech talk. It will take some time for me to get into more detail of important things about cache, cache server, leader and follower, data center, replicate data center, read availability those interesting topics. 

My notes

I like to take some time to write down the notes on the paper, so that I can train myself in terms of high level design.  



Cenovus, Suncor choose different paths on restoring oil output as prices recover

 Here is the article. 

What is the difference? Suncor shutdown of one of two production trains while Cenovus brought back 60,000 barrels a day in June. 

CALGARY — Two of Canada's biggest oilsands producers are taking different approaches to restoring production that was shut down during the pandemic lockdowns as signs of an economic recovery and higher oil prices emerge.

Both Cenovus Energy Inc. and Suncor Energy Inc. reported millions of dollars in losses in the second quarter as they throttled back oil production amid a global crude market awash in surplus barrels.

However, while Cenovus brought back 60,000 barrels a day in June to take advantage of higher bitumen prices, Suncor signalled Thursday the shutdown of one of the two production trains at its 194,000-bpd Fort Hills oilsands mine could remain in place for some time.

"In response to the sharp decline in oil prices in April, we quickly reduced production volumes at our oilsands operations while continuing to steam and store the mobilized oil in the reservoir," said Cenovus CEO Alex Pourbaix on a Thursday conference call.

"When the market price for Western Canadian Select (bitumen-blend oil) increased almost 10-fold in June compared with April, we acted fast to ramp our oilsands production back up to take advantage of the improved pricing."

The increase in WCS to an average of C$46.03 per barrel in June from C$4.92 in April prompted a boost in production to record levels at Cenovus's Christina Lake operations in northeastern Alberta, the Calgary-based company said.

System design: TAO - Improve availability - read failover

 August 11, 2020

Introduction

It takes a lot of wisdom to learn system design in less than two weeks. My goal is to start to build strong interest to learn large distributed system, no matter what if I can get Facebook offer or not. I like to learn one system design, TAO. 

TAO - improve availability - read failover

There are so many ways to learn. I spent last two days to study TAO, play the same video over 10 times, and also read the paper, and then read slides, and ask myself question. 

I also start to draw the system design high level by myself, using my own pen and paper. 



More handwriting notes



2 TSX Energy Stocks with Dividend Yields of Up to 6.8%

 Here is the article.

Quick highlights from the article.

  1. CNQ - forward yield 6.8% - what is forward?
  2. CNQ - 50% revenue drop from the last year, $310 million loss compared to $2.8 billion profit
  3. CNQ - cash flow - $415 million in Q2 of 2020
  4. CNQ - 5 dolloar a share gain analyst expected

During times of volatility, investors look for a steady stream of income. The easiest way to do that is to research stocks that pay out a high dividend. The energy and power sectors are good segments to hunt for stocks with good yields. Here we look at two such as TSX giants that should interest dividend investors.

Canadian Natural Resources has a forward yield of 6.8% 

When your revenue drops by 50% from the last year and your report a$310 million loss compared to a profit of $2.8 billion in the previous year and you still beat analyst expectations, the market knows you’re a force to be reckoned with.

Canadian Natural Resources (TSX:CNQ)(NYSE:CNQ) (TSX: CNQ) is a giant in the fossil fuels segment. The company reported an adjusted cash flow of $415 million in Q2 of 2020. It produced 1.16 million barrels of oil equivalent per day, including about 922,000 barrels per day of crude and natural gas liquids, up from 1.02 million boe/d including 770,000 bbl/d in 2019.

The world is slowly transitioning away from fossil fuels to renewable sources of energy but you can bet your last buck that CNQ will find a way to go about its business. As oil prices revert to pre-pandemic levels, CNQ’s cash positions will only get stronger and we can expect a return to profitability.

The stock is currently trading at $26.37 and analysts have given it an average price target of $31.27, an upside of over 18%. Combine this with a sweet dividend yield of 6.79%, and the savvy investor could be sitting on a pretty profit in 12 months.

A renewable growth stock

Transalta Renewables (TSX:RNW) is one of the largest renewable power producers in Canada with 2,527 MW of power across North America and Australia. The company has interests in the wind, hydro, natural gas, and solar segments.

Wind energy is a major driver for Transalta, generating 51% of cash flows followed by natural gas, hydro, and solar at 43%, 4%, and 2%, respectively. The company has been an early mover in the renewable energy space

It reported its numbers for the second quarter of 2020 recently. Its EBITDA was $115 million, a $4 million increase to the same period in 2019, no small achievement during the pandemic.

Adjusted funds from operations (AFFO) were $90 million, a $10 million, or 13% increase to the second quarter in 2019. Transalta states that it pays out “80 to 85% of cash available for distribution to the shareholders of the company on an annual basis.” This bodes well for investors in Transalta. The company sports a strong dividend yield of 6.08% which should hold investors in good stead during uncertain times.

The stock is currently trading at $15.78 and is poised to reap benefits in the long term as the world transitions from fossil fuels to renewable sources for its energy requirements. I had written about the stock in April this year when it was trading at $13.86 and recommended a buy at that time. My view is still the same: the company has tremendous growth potential.

Both stocks operate in the power space and if given a choice, I would go with Transalta. It is the stock of the future.

The post 2 TSX Energy Stocks with Dividend Yields of Up to 6.8% appeared first on The Motley Fool Canada.



AC.TO research: Air Canada to launch revamped Aeroplan amid devastated travel industry

 Here is the link. 


Wikipedia: Thundering herd problem

 Here is the article. 

In computer science, the thundering herd problem occurs when a large number of processes or threads waiting for an event are awoken when that event occurs, but only one process is able to handle the event. When the processes wake up, they will each try to handle the event, but only one will win. All processes will compete for resources, possibly freezing the computer, until the herd is calmed down again.[1]

Mitigation[edit]

The Linux-kernel will serialize responses for requests to a single file descriptor, so only one thread (process) is woken up.[2]

Similarly in Microsoft Windows, I/O completion ports can mitigate the thundering herd problem, as they can be configured such that only one of the threads waiting on the completion port is woken up when an event occurs.[3]

In systems which rely on a backoff mechanism (e.g. exponential backoff), the clients will retry failed calls, by waiting a specific amount of time between consecutive retries. In order to avoid the thundering herd problem, jitter can be purposefully introduced, in order to break the synchronization across the clients thereby avoiding collisions. In this approach, randomness is added to the wait intervals between retries, so that clients are no longer synchronized.


Actionable Items

I  tried to understand Facebook TAO data store, and thundering herd problem is one of important concepts for me to learn in order to understand the requirements. 





Monday, August 10, 2020

SPG stock: The biggest US mall owner Simon says still looking to salvage other distressed retailers

 Here is the article. 


Suncor stock: My quick sharing with a friend

 August 10, 2020

Introduction

I like to share my experience with SU.TO stock purchase history with a friend quickly. I do think that it is better to invest on IMO.TO stock since it can go up more and more volatile. 

SU.TO stock


Comparison to IMO.TO


Comparison to CNQ.TO

Comparison to OVV.TO


THC stock: My purchase of 10 shares - 20% gains from June 8, 2020

 August 10, 2020


An Introduction to Look-Aside Caching

 Here is the article. 

Performance is critical to the success of any given microservice.  Overall performance is the result of applying ‘performance friendly’ techniques at various points in the design, development, and delivery of microservices. In many cases, however, you can make vast performance improvements through basic techniques like implementing and optimizing caching at various points in between the consumers of data (users and applications) and servers that store data. Caches can return data much faster than the disk-based databases that originate the data because of caches’ use of memory to provide lower latency access. Caches are also usually located much closer to the consumers of data from a network topology perspective.

A cache can be inserted anywhere in the infrastructure where there is congestion with data delivery. In this post, we’ll focus on look-aside caching that serves as a highly performant alternative to accessing data from a microservice’s backing store. We will also clarify the meaning of various terms associated with caching patterns - such as read-aside, read thru, write through, and write behind caches - and when to choose each pattern.

Look-Aside Cache vs. Inline Cache

The two main caching patterns are the look-aside caching pattern and the inline caching pattern. The descriptions and differences between these patterns are shown in the table below.

Look-aside cache

How it reads - explain it to myself

Application requests data from cache

Cache delivers data, if available

If data not available, application gets data from backing store and writes it to the cache for future requests (read aside)

How it writes - explain to myself

Application writes new data or updates to existing data in both the cache and the backing store - or - all writes are done to the backing store and the cache copy is invalidated. 

Pattern - look-aside cache

Inline cache - pattern, try to separate from look-aside cache

How it reads - explain it to myself


How it writes - explain to myself

Vacation days: August 10 to August 13

 August 10, 2020

Introduction

It is my four days vacation. I like to make some plans and also really enjoy learning in those four days. 

Four day vacation

I decided to take four days vacation and stay at home, I need to study and really prepare myself for important things in my life. System design learning and algorithm practice. 



BIGC stock: BigCommerce Announces Closing of Initial Public Offering and Exercise in Full of the Underwriters’ Option to Purchase Additional Shares

 Here is the article. 


BIGC stock: New IPO e-commerce company

 August 10, 2020

Introduction

It is not an easy task to work on stock market investment. Sometimes I need a few people to bounce back the ideas, so I started this investment wechat group. We talk and share ideas how to explore the stock market together. 

BIGC stock

My friend shared with me that she did a few transactions in the morning, get in and get out from SAVE stock; next she went into BIGC stock. 



Why Alteryx Stock Plunged Today

 Here is the link. 

Fast growth, volatile stock

Data analytics companies garnered a lot of attention in 2019. Google parent Alphabet purchased Looker for $2.6 billion, and salesforce.com bought industry leader Tableau for $15.7 billion. And those left over notched massive growth. Alteryx specifically ended 2019 having brought in $418 million in revenue, a 65% increase. Pretty impressive, especially considering it built on a 55% revenue gain in 2018.

In the early days of the pandemic, Alteryx started the new decade strong. Sales grew by a massive amount again year over year, and the software platform continued operating at an exceptionally high gross profit margin (revenue less cost of the service). The dollar-based net expansion rate was 128%, implying that in addition to the company signing up new customers, existing ones spent an average 28% more with Alteryx than they did a year ago.  

Metric

Full Year 2019

Full Year 2018

Change

Revenue

$109 million

$76.0 million

43%

Adjusted gross profit margin

91.3%

90.5%

(0.8 pp)

Adjusted net income (loss)

($6.45 million)

$2.99 million

N/A

Free cash flow

$15.0 million

$14.5 million

3%

Key largo: How to actively make adjustment in short term?

 August 10, 2020

Introduction

It is interesting to learn from peers in my wechat investment group. The idea to invest is different from passive investment. I like to write down and think more carefully later on. 


Key largo US retirement account

I got into stocks position from ETF starting from June 8, 2020 with big purchase of $13,000 US dollars, now the balance is $3000 dollars less compared to the peak $22,000. How to make adjustment on market swing?

Today SAVE went up above 8%. The idea is to sell INTL and WDC stock, and then purchase SAVE stock. Try to catch gains first, and then purchase back INTL and WDC stock. 

Right now I have 100% equities. 






AC stock: My investment as a beginner

 August 10, 2020

Introduction

It is my short term investment. I like to write down lessons I learn from market swing. This is the first time I made some profit from Air Canada stock. 

My gains



MGM stock: My story as a beginner

 August 10, 2020

Introduction

As a beginner, I spent time to look into MGM stock and business analysis. I do think that learning experience is so interesting to write down and share. It is not long ago that I did purchase 300 shares, at that time, I worried about loss of my Key largo portfolio. So it is better to share the detail. 

My story

I did purchase 300 shares of MGM, and then sold it in a hurry. At that time, I think that SPX 500 index went down around 3000. 

Right now, if I can hold those stocks of MGM, then I will have 1500 dollars gains easily. As a beginner, I do enjoy learning process, I knew how MGM works out it's debt, furlough the employees, and it should not have problems in short term. 

It is a safe bet. But it is also perfect investment for short term return. 

MGM: Diller’s IAC/InterActive Takes 12% Stake in MGM for $1 Billion

 Here is the article. 

Here are highlights:

  1. MGM shares jumped 22%
  2. Online gaming market represents a $450 billion global opportunity
  3. IAC InterActive Corp. said it has built a 12% stake in hospitality and entertainment company MGM Resorts International for about $1 billion.

(Bloomberg) -- IAC InterActive Corp. said it has built a 12% stake in hospitality and entertainment company MGM Resorts International for about $1 billion.

“With the separation of Match Group from IAC, and ‘new’ IAC emerging with $3.9 billion of cash, no debt, and its opportunistic zeal intact, we are energized and excited to make this investment in MGM,” IAC Chairman Barry Diller said in a statement Monday.

MGM shares jumped 22% and IAC shares rose 1.4% Monday morning in New York after the announcement.

One thing that attracted Diller to MGM in particular is an area that currently comprises a tiny portion of its revenue – online gaming. That market represents a $450 billion global opportunity, according to IAC, with less than 10% penetration online.

In a letter to shareholders, Diller said investors might be surprised by the move. It’s unusual for IAC to purchase a large stake in a public company, he said. Furthermore, IAC purchased MGM’s common equity, “the exact same securities that any investor with exactly $19 could buy and sell any day in the market.”

IAC is also buying securities in a business that has relatively little to do with the internet today, veering from its traditional strategy. The reason is “that we believe MGM presented a ‘once in a decade’ opportunity for IAC to own a meaningful piece of a preeminent brand in a large category with great potential to move online.”

MGM, like other casino operators, has been hit hard by the coronavirus, which triggered a monthslong closing of its properties in the U.S. and a severe contraction in Macau.

The company is in a position to weather the storm, having sold nearly all of its resorts to investors in a sale-leaseback arrangement that freed up billions in cash. Still, MGM has cut staff and furloughed others as it copes with far less business due to the virus.

The company last month gave its chief executive officer position permanently to Bill Hornbuckle, a company veteran who had been acting CEO since March. In a previous role, as marketing chief, Hornbuckle spearheaded MGM’s customer-loyalty program, which IAC cited as one of the enticing aspects of the company.



HNGR: A Look At The Intrinsic Value Of Hanger, Inc. (NYSE:HNGR)

 Here is the link. 


Sunday, August 9, 2020

USENIG.org: TAO: Facebook’s Distributed Data Store for the Social Graph

August 9, 2020

Introduction

It is the first time I came cross the original conference paper. I like to read the paper as well. 

TAO: Facebook’s Distributed Data Store for the Social Graph

Here is the link of paper. Total pages are 12 pages. I like to spend 30 minutes to 60 minutes to read the paper. 

Statement: 

TAO, a read-optimized graph data store we have built to handle a demanding Facebook workload

read-optimized graph data store - 

Web server -> TAO data store -> Database, memcache as a lookaside cache - what is lookaside cache

MySQL - persistent storage - graph-aware cache - how to define graph-aware? 

CAP - AP - favors availability and per-machine efficiency over strong consistency. 

TAO can sustain a billion reads per second on a changing data set of many petabytes. 

efficient and available read-mostly access to a changing graph 

objects and associations, a data model and API that we use to access the graph

TAO - a distributed system that implements this API 

Facebook has more than a billion active users who record their relationships, share their interests, upload text, images, and video, and curate semantic information about their data [2]. The personalized experience of social applications comes from timely, efficient, and scalable access to this flood of data, the social graph. In this paper we introduce TAO, a read-optimized graph data store we have built to handle a demanding Facebook workload. Before TAO, Facebook’s web servers directly accessed MySQL to read or write the social graph, aggressively using memcache [21] as a lookaside cache. TAO implements a graph abstraction directly, allowing it to avoid some of the fundamental shortcomings of a lookaside cache architecture. TAO continues to use MySQL for persistent storage, but mediates access to the database and uses its own graph-aware cache. TAO is deployed at Facebook as a single geographically distributed instance. It has a minimal API and explicitly favors availability and per-machine efficiency over strong consistency; its novelty is its scale: TAO can sustain a billion reads per second on a changing data set of many petabytes. Overall, this paper makes three contributions. We motivate (§ 2) and characterize (§ 7) a challenging workload: efficient and available read-mostly access to a changing graph. We describe objects and associations, a data model and API that we use to access the graph (§ 3). Lastly, we detail TAO, a geographically distributed system that implements this API (§§ 4–6), and evaluate its performance on our workload (§ 8).

Let me try to read the 2.1 Serving the graph from Memcache

Facebook was originally built by storing the social graph in MySQL, querying it from PHP, and caching results in memcache [21]. This lookaside cache architecture is well suited to Facebook’s rapid iteration cycles, since all of the data mapping and cache-invalidation computations are in client code that is deployed frequently. Over time a PHP abstraction was developed that allowed developers to read and write the objects (nodes) and associations (edges) in the graph, and direct access to MySQL was deprecated for data types that fit the model. TAO is a service we constructed that directly implements the objects and associations model. We were motivated by encapsulation failures in the PHP API, by the opportunity to access the graph easily from non-PHP services, and by several fundamental problems with the lookaside cache architecture:

TAO is a service - objects and associations model 

Inefficient edge lists: A key-value cache is not a good semantic fit for lists of edges; queries must always fetch the entire edge list and changes to a single edge require the entire list to be reloaded. Basic list support in a lookaside cache would only address the first problem; something much more complicated is required to coordinate concurrent incremental updates to cached lists.

coordinate concurrent incremental updates to cached lists?  - how to coordinate? 

Distributed control logic: In a lookaside cache architecture the control logic is run on clients that don’t communicate with each other. This increases the number of failure modes, and makes it difficult to avoid thundering herds. Nishtala et al. provide an in-depth discussion of the problems and present leases, a general solution [21]. For objects and associations the fixed API allows us to move the control logic into the cache itself, where the problem can be solved more efficiently.

What is thundering herds? - 

Expensive read-after-write consistency: Facebook uses asynchronous master/slave replication for MySQL, which poses a problem for caches in data centers using a replica. Writes are forwarded to the master, but some time will elapse before they are reflected in the local replica. Nishtala et al.’s remote markers [21] track keys that are known to be stale, forwarding reads for those keys to the master region. By restricting the data model to objects and associations we can update the replica’s cache at write time, then use graph semantics to interpret cache maintenance messages from concurrent updates. This provides (in the absence of multiple failures) read-after-write consistency for all clients that share a cache, without requiring inter-regional communication.

asynchronous master/ slave replication for MySQL - caches in data centers using a replica - 

Writes are forwarded to the master, but some time ... 

update replica's cache at write time, then ...

read-after-write consistency - 


2.2 TAO’s Goal TAO provides basic access to the nodes and edges of a constantly changing graph in data centers across multiple regions. It is optimized heavily for reads, and explicitly favors efficiency and availability over consistency. A system like TAO is likely to be useful for any application domain that needs to efficiently generate fine grained customized content from highly interconnected data. The application should not expect the data to be stale in the common case, but should be able to tolerate it. Many social networks fit in this category


3 TAO Data Model and API 

Facebook focuses on people, actions, and relationships. We model these entities and connections as nodes and edges in a graph. This representation is very flexible; it directly models real-life objects, and can also be used to store an application’s internal implementation-specific data. TAO’s goal is not to support a complete set of graph queries, but to provide sufficient expressiveness to handle most application needs while allowing a scalable and efficient implementation. 


Consider the social networking example in Figure 1a, in which Alice used her mobile phone to record her visit to a famous landmark with Bob. She ‘checked in’ to the Golden Gate Bridge and ‘tagged’ Bob to indicate that he is with her. Cathy added a comment that David has ‘liked.’ The social graph includes the users (Alice, Bob, Cathy, and David), their relationships, their actions (checking in, commenting, and liking), and a physical location (the Golden Gate Bridge). Facebook’s application servers would query this event’s underlying nodes and edges every time it is rendered. Fine-grained privacy controls mean that each user may see a different view of the checkin: the individual nodes and edges that encode the activity can be reused for all of these views, but the aggregated content and the results of privacy checks cannot.

Social graph - users - Alice, Bob, Cathy, and David

their relationships, their actions - checking in, commenting, and liking

a physical location (the Golden Gate Bridge)

3.1 Objects and Associations

TAO objects are typed nodes, and TAO associations are typed directed edges between objects. Objects are identified by a 64-bit integer (id) that is unique across all objects, regardless of object type (otype). Associations are identified by the source object (id1), association type (atype) and destination object (id2). At most one association of a given type can exist between any two objects. Both objects and associations may contain data as key→value pairs. A per-type schema lists the possible keys, the value type, and a default value. Each association has a 32-bit time field, which plays a central role in queries1. 

Object: (id) → (otype, (key  value)∗) 

Assoc.: (id1, atype, id2) → (time, (key  value)∗) 

Figure 1b shows how TAO objects and associations might encode the example, with some data and times omitted for clarity. The example’s users are represented by objects, as are the checkin, the landmark, and Cathy’s comment. Associations capture the users’ friendships, authorship of the checkin and comment, and the binding between the checkin and its location and comments.

Actions may be encoded either as objects or associations. Both Cathy’s comment and David’s ‘like’ represent actions taken by a user, but only the comment results in a new object. Associations naturally model actions that can happen at most once or record state transitions, such as the acceptance of an event invitation, while repeatable actions are better represented as objects. 

Although associations are directed, it is common for an association to be tightly coupled with an inverse edge. In this example all of the associations have an inverse except for the link of type COMMENT. No inverse edge is required here since the application does not traverse from the comment to the CHECKIN object. Once the checkin’s id is known, rendering Figure 1a only requires traversing outbound associations. Discovering the checkin object, however, requires the inbound edges or that an id is stored in another Facebook system. 

The schemas for object and association types describe only the data contained in instances. They do not impose any restrictions on the edge types that can connect to a particular node type, or the node types that can terminate an edge type. The same atype is used to represent authorship of the checkin object and the comment object in Figure 1, for example. Self-edges are allowed.

3.2 Object API 

TAO’s object API provides operations to allocate a new object and id, and to retrieve, update, or delete the object associated with an id. A notable omission is a compare-and-set functionality, whose usefulness is substantially reduced by TAO’s eventual consistency semantics. The update operation can be applied to a subset of the fields. 

3.3 Association API 

Many edges in the social graph are bidirectional, either symmetrically like the example’s FRIEND relationship or asymmetrically like AUTHORED and AUTHORED BY. Bidirectional edges are modeled as two separate associations. TAO provides support for keeping associations in sync with their inverses, by allowing association types to be configured with an inverse type. For such associations, creations, updates, and deletions are automatically coupled with an operation on the inverse association. Symmetric bidirectional types are their own inverses. The association write operations are: 

• assoc add(id1, atype, id2, time, (k→v)*) – Adds or overwrites the association (id1, atype,id2), and its inverse (id1, inv(atype), id2) if defined. 

• assoc delete(id1, atype, id2) – Deletes the association (id1, atype, id2) and the inverse if it exists. 

• assoc change type(id1, atype, id2, newtype) – Changes the association (id1, atype, id2) to (id1, newtype, id2), if (id1, atype, id2) exists.


3.4 Association Query API

The starting point for any TAO association query is an originating object and an association type. This is the natural result of searching for a specific type of information about a particular object. Consider the example in Figure 1. In order to display the CHECKIN object, the application needs to enumerate all tagged users and the most recently added comments. 

A characteristic of the social graph is that most of the data is old, but many of the queries are for the newest subset. This creation-time locality arises whenever an application focuses on recent items. If the Alice in Figure 1 is a famous celebrity then there might be thousands of comments attached to her checkin, but only the most recent ones will be rendered by default. 


TAO’s association queries are organized around association lists. We define an association list to be the list of all associations with a particular id1 and atype, arranged in descending order by the time field: 

Association List: (id1, atype) → [anew ...aold] 

For example, the list (i, COMMENT) has edges to the example’s comments about i, most recent first. 

TAO’s queries on associations lists: 

• assoc get(id1, atype, id2set, high?, low?) – returns all of the associations (id1, atype, id2) and their time and data, where id2 ∈ id2set and high ≥ time ≥ low (if specified). The optional time bounds are to improve cacheability for large association lists (see § 5). 

• assoc count(id1, atype) – returns the size of the association list for (id1, atype), which is the number of edges of type atype that originate at id1. 

• assoc range(id1, atype, pos, limit) – returns elements of the (id1, atype) association list with index i ∈ [pos,pos+limit). 

• assoc time range(id1, atype, high, low, limit) – returns elements from the (id1, atype) association list, starting with the first association where time ≤ high, returning only edges where time ≥ low. 


TAO enforces a per-atype upper bound (typically 6,000) on the actual limit used for an association query. To enumerate the elements of a longer association list the client must issue multiple queries, using pos or high to specify a starting point. 

For the example shown in Figure 1 we can map some possible queries to the TAO API as follows: 

• “50 most recent comments on Alice’s checkin” ⇒ assoc range(632, COMMENT, 0, 50) • “How many checkins at the GG Bridge?” ⇒ assoc count(534, CHECKIN)

4 TAO Architecture 

In this section we describe the units that make up TAO, and the multiple layers of aggregation that allow it to scale across data centers and geographic regions. TAO is separated into two caching layers and a storage layer. 

4.1 Storage Layer 

Objects and associations were stored in MySQL at Facebook even before TAO was built; it was the backing store for the original PHP implementation of the API. This made it the natural choice for TAO’s persistent storage. 

The TAO API is mapped to a small set of simple SQL queries, but it could also be mapped efficiently to range scans in a non-SQL data storage system such as LevelDB [3] by explicitly maintaining the required indexes. When evaluating the suitability of a backing store for TAO, however, it is important to consider the data accesses that don’t use the API. These include backups, bulk import and deletion of data, bulk migrations from one data format to another, replica creation, asynchronous replication, consistency monitoring tools, and operational debugging. An alternate store would also have to provide atomic write transactions, efficient granular writes, and few latency outliers. 

Database - backups, bulk import and deletion of data, bulk migration from one data format to another, replica creation, asynchronous replication, consistency monitoring tools, and operational debugging. 



Given that TAO needs to handle a far larger volume of data than can be stored on a single MySQL server, we divide data into logical shards. Each shard is contained in a logical database. Database servers are responsible for one or more shards. In practice, the number of shards far exceeds the number of servers; we tune the shard to server mapping to balance load across different hosts. By default all object types are stored in one table, and all association types in another. 

Each object id contains an embedded shard id that identifies its hosting shard. Objects are bound to a shard for their entire lifetime. An association is stored on the shard of its id1, so that every association query can be served from a single server. Two ids are unlikely to map to the same server unless they were explicitly colocated at creation time.

Actionable Items

I finished reading at 11:23 PM. So it took me more than two hours. 

I do think that it is better for me to understand the paper; Best way is to listen the presentation again after paper reading. 



My cheat plan: Watch TAO talk over 10 times - 200 minutes

 August 9, 2020

Introduction

It is a small project I like to work on. How to build strong foundation of system design? What I like to do is to play TAO talk 10 times, and each time I like to find some articles to read and then I can understand the topic better. 

10 times repetitions

The idea I can look into can be the following:

  1. Sharding in cache vs sharding in database
  2. Lead cache and follower cache - 
  3. Write-through cache topic
  4. TAO - object associations - product engineer - a small chat about my understanding
  5. Challenge - high read availability - how to work on the challenge?
  6. Failover and replication data center
  7. ...

Actionable items


I did spend almost 10 times to watch the videos, and I do think that it is so helpful for me to learn one system design very well. So I can use it for the template in short future. 

This TAO design covers a few good topics, scalable, high availability of reading, timelyness of reading, and social graph, how to model social graph using TAO - association and object, and how to design cache considering data center and leader and follower caches. 

Follow up

August 15, 2020

TAO tech talk link is here. 

Write-Through Cache

 Write-through cache is a caching technique in which data is simultaneously copied to higher level caches, backing storage or memory. It is common in processor architectures that perform a write operation on cache and backing stores at the same time.

Write-through cache helps increase read performance in memory access methods because the requested data is already present in the cache and memory. In each write-through operation, data that is brought into the cache is also written into the backing store, which is the primary memory (RAM, in most cases).

Write-through cache also helps with data recovery, as the data in operation is written to both cache and memory. It is nearly impossible to restore data from cache. However, to an extent, modern operating systems (OS) have the ability to save an instance of running memory.

Facebook tech talk: TAO: Facebook’s Distributed Data Store for the Social Graph

 Here is the link. 

Abstract: 

We introduce a simple data model and API tailored for serving the social graph, and TAO, an implementation of this model. TAO is a geographically distributed data store that provides efficient and timely access to the social graph for Facebook’s demanding workload using a fixed set of queries. It is deployed at Facebook, replacing memcache for many data types that fit its model. The system runs on thousands of machines, is widely distributed, and provides access to many petabytes of data. TAO can process a billion reads and millions of writes each second.


Tao summary

Efficiency at scale 

Read latency 

The solution is to separate cache and DB, graph-specific caching, subdivide data centers

Write timeliness

The solution is to write-through cache, asynchronous replication

Read availability 

The solution is to alternate data sources


More details:

web server -> query graph -> html 

when to aggregate and rendering ...

query them and dynamically - hard to predict, different from each user, each user has different view of graph

guarantee for right content - TAO 

database, stacks, dynamic rendering every time - highest ...

query in TAO - multiple billions per second - whole graph - petabytes data every data centers

Limitation 

Scale - large level 

efficiency in scale - not too much cost 

Dynamic resolution of data dependencies - post - three rounds - HTTP 

Low read latency - read from local data center 

Timeliness of writes - data center - web server - read - cannot read after write - high readability - TAO 

Graph in Memcache 

PHP abstraction - obj & Assoc API - memcache (nodes, edges, edge lists)

Service -> control - cache - model 

objects= Nodes 

64-bit IDs type, with a schema for fields 

Associations - Edges 

Association lists 

ide1, type - descending order by time 

Every edge has time - all queries more than one edge - size of association lists 

objects and associations API - read work very well

web servers - stateless 

cache - objects, assoc lists, assoc counts, Database TAO ,

stateless, shared by id, servers -> read qps 

Sharded by id - servers - bytes

subdivding the data center  

web servers, cache, database, 

servers - nearby building - many open sockets, lots of hot spot 

Subdividing the data center 

web servers - cache - database 

distributed write control logic - cache - thundering objects 

leader cache - database - follower cache

Timeliness of writes 


Async DB replication 

Master data center   Replica data center

web servers 

writes forwarded to master

Inval and forwarded to master 

Inval and refill embedded in SQL 


Improving availability: Read failover 

TAO summary 

Efficiency at scale / Read latency 

Write timeliness - write-through cache, asynchronous replication 

Read availability - Alternate data sources

Questions: 

MySQL - 

graphic roles - 

clean separations - 

size of Node - largest 1 MB - graph, long tail large value - content link, not a photo itself, metadata, ...

17,000 lines of comment - TAO ... stretch to limit 

Iphone - ...

Leader node - half hour ...... 

leaders do - reduce temperature of hot spots - 

TAO - not designed heavy write - less timely, more data intensity 

who to suggest for friends

Consistency model - TAO - opportunity - 

TAO summary - efficiency at scale / Read latency / Write timeliness/ Read availability