đź’» The byte is not enough
291 subscribers
1 photo
10 files
64 links
Personal hand-picked collection of articles and tutorials on the matter of software engineering and computer science. Occasionally on science, history, or linguistics.

@virtyaluk for any inquiries.

https://modern-dev.com/
https://github.com/virtyaluk
Download Telegram
nsdi13-final170_update.pdf
370.1 KB
🧑‍🔬 Scaling Memcache at Facebook

While reading Facebook's Twine cluster management system white paper, I noticed an interesting thing in section 5 where paper authors claim to have a highly optimized memcached deployment that can handle "930k lookups per second on an 18-core/36-hyperthread machine". Just think for a second, 930 000 lookup requests per second on a single machine. How the heck they could achieve this kind of performance. To answer this and many other questions, I went on reading another Facebook white paper on optimizing Memcached for meeting the world's largest social network needs.

A short version listing all the big things Facebook incorporated into their version of Memcached can be found here:

https://medium.com/@shagun/scaling-memcache-at-facebook-1ba77d71c082

#systemdesign #facebook #memcached #distributed #caching #design #reliability #performance #scalability #faulttolerance #research #paper #redis #dht

đź’» The byte is not enough
atc13-bronson.pdf
729.7 KB
🧑‍🔬 TAO: Facebook’s Distributed Data Store for the Social Graph

Have you ever wondered how Facebook manages its social graph, containing petabytes of user data? What techniques do they apply to serve billions of reads and millions of writes each second? All this and much more in another great white paper on Facebook's TAO - 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.

https://engineering.fb.com/2013/06/25/core-data/tao-the-power-of-the-graph/

USENIX ATC '13 - TAO: Facebook’s Distributed Data Store for the Social Graph:
https://www.youtube.com/watch?v=sNIvHttFjdI

#systemdesign #facebook #memcached #tao #socialgraph #api #caching #design #reliability #performance #scalability #faulttolerance #consistency #research #paper #mysql #infrastructure

đź’» The byte is not enough
paxos-simple.pdf
92.8 KB
🧑‍🔬 Paxos Made Simple

In the year 1989 Leslie Lamport, a known computer scientist in the field of distributed systems, published his tremendous work on Paxos — a family of protocols for solving consensus in a network of unreliable or fallible processors. Though from the very beginning, the proposed algorithm was diminished by computer science society due to its complexity, it started gaining significant recognition after almost 10 years since first published having a second coming in 1998. In late 2001, Lamport published a simplified version of the original paper, discarding unnecessary information and providing the description on the backbone of Paxos protocol.

Paxos Simplified by Chris Colohan:
https://www.youtube.com/watch?v=SRsK-ZXTeZ0

#systemdesign #paxos #consensus #lamport #design #reliability #performance #faulttolerance #scalability #consistency #quorum #research #paper #infrastructure #distributed

đź’» The byte is not enough
16cb30b4b92fd4989b8619a61752a2387c6dd474.pdf
186.2 KB
🧑‍🔬 MapReduce: Simplified Data Processing on Large Clusters

Another seminal work from the past that established distributed systems' evolution for decades ahead. The MapReduce model is probably the most well know programming model designed for processing and generating big data sets with a parallel, distributed algorithm on a cluster.

#systemdesign #mapreduce #design #performance #faulttolerance #scalability #research #paper #infrastructure #distributed #computation #cluster #gfs

đź’» The byte is not enough
Turbine_Facebook’s_Service_Management_Platform_for_Stream_Processing.pdf
1.9 MB
🧑‍🔬 Turbine: Facebook’s Service Management Platform
for Stream Processing


A scalable service management platform for Facebook’s stream processing service. Turbine is designed to bridge the gap between the capabilities of existing general-purpose cluster management frameworks like Tupperware and Facebook’s stream processing requirements. In production for several years now, Turbine has enabled a boom in stream processing at Facebook.

https://engineering.fb.com/2020/04/21/data-infrastructure/turbine/

#systemdesign #turbine #facebook #streaming #processing #design #performance #acid #reliability #faulttolerance #scalability #research #paper #infrastructure #distributed #cluster #management

đź’» The byte is not enough
👍1
LogDevice: a distributed data store for logs

A log is the simplest way to record an ordered sequence of immutable records and store them reliably. Build a data intensive distributed service and chances are you will need a log or two somewhere. At Facebook, we build a lot of big distributed services that store and process data. Want to connect two stages of a data processing pipeline without having to worry about flow control or data loss? Have one stage write into a log and the other read from it. Maintaining an index on a large distributed database? Have the indexing service read the update log to apply all the changes in the right order. Got a sequence of work items to be executed in a specific order a week later? Write them into a log, have the consumer lag a week. Dream of distributed transactions? A log with enough capacity to order all your writes makes them possible. Durability concerns? Use a write-ahead log.

https://engineering.fb.com/2017/08/31/core-data/logdevice-a-distributed-data-store-for-logs/

#systemdesign #logdevice #log #wal #logsdb #rocksdb #consensus #paxos #quorum #lsmtree #facebook #storage #design #performance #scalability #research #faulttolerance #infrastructure

đź’» The byte is not enough
Database Storage Engines: B-Tree vs LSM-Tree

Have you ever concern yourself with the question of how modern database systems' internal storage works? Well, there are two popular ways to handle data storage — a B-Tree (a generalization of Binary Search Tree) and a Log-Structured Merge Tree, both with having pros and cons.

Here are a few short articles to get a grasp of the trade-offs between the two:

1. https://blog.yugabyte.com/a-busy-developers-guide-to-database-storage-engines-the-basics/
2. https://blog.yugabyte.com/a-busy-developers-guide-to-database-storage-engines-advanced-topics

3. https://rkenmi.com/posts/b-trees-vs-lsm-trees

4. https://tikv.org/deep-dive/key-value-engine/b-tree-vs-lsm/

#systemdesign #dbs #databases #design #btree #lsmtree #storage #performance #sql #nosql #engine

đź’» The byte is not enough
zab.totally-ordered-broadcast-protocol.2008.pdf
264.7 KB
🧑‍🔬 Architecture of ZAB – ZooKeeper Atomic Broadcast protocol

The ZAB protocol ensures that the Zookeeper replication is done in order and is also responsible for the election of leader nodes and the restoration of any failed nodes. In a Zookeeper ecosystem, the leader node is the heart of everything; every cluster has one leader node and the rest of the nodes are followers. All incoming client requests and state changes are received at first by the leader with responsibility to replicate it across all its followers (and itself). All incoming read requests are also load balanced by the leader within itself and its followers.

Original ZAB paper:
https://marcoserafini.github.io/papers/zab.pdf

Implementation details of ZAB:
http://www.tcs.hut.fi/Studies/T-79.5001/reports/2012-deSouzaMedeiros.pdf

#systemdesign #zab #consensus #total #order #broadcast #design #quorum #performance #faulttolerance #yahoo #paper #research #scalability

đź’» The byte is not enough
Building Facebook’s service encryption infrastructure

What does it take to incorporate security protocols into a system of millions of services running on thousands of machines across the world? You could use some well-known technologies as Kerberos, but will soon realize it does not work on a big scale. You would then probably stick to the idea of building a custom in-house solution, which is exactly what Facebook did to suit their needs in security and operability without compromising performance.

https://engineering.fb.com/2019/05/29/security/service-encryption/

#systemdesign #security #facebook #kerberos #tls #encryption #traffic #networking #dns #performance #faulttolerance #scalability #infrastructure

đź’» The byte is not enough
🧑‍🔬 CORFU: A Distributed Shared Log

Despite almost forty years of research into replicated storage schemes, the only approach so far to scale up capacity and throughput has been to shard data and trade consistency for performance. The CORFU system breaks this seeming tradeoff by organizing a cluster of drives as a single, shared log. CORFU offers a single-copy semantics at cluster-scale speeds, providing a scalable source of atomicity and durability for distributed systems.

https://www.youtube.com/watch?v=GmVQVT9aZfU

Original paper:
https://www.cs.utexas.edu/~lorenzo/corsi/cs380d/papers/a10-balakrishnan.pdf

#systemdesign #corfu #microsoft #design #reliability #algorithms #distributed #log #replication #consensus #performance #scalability #faulttolerance #paper #research #infrastructure

đź’» The byte is not enough
🧑‍🔬 Delos: Simple, flexible storage for the Facebook control plane

Facebook Delos — is a fundamentally new architecture for building replicated storage systems. Its modular, layered design provides flexibility and simplicity without sacrificing performance or reliability. Delos enables fast time-to-deployment for new storage systems — Facebook deployed an initial version in production within eight months — as well as safe, rapid evolution. Facebook swapped in a new ordering mechanism to obtain 10x lower latency without any service downtime.

https://www.youtube.com/watch?v=wd-GC_XhA2g&t=314s

Blog post presentation:
https://engineering.fb.com/2019/06/06/data-center-engineering/delos/

Original paper:
https://www.usenix.org/system/files/osdi20-balakrishnan.pdf

#systemdesign #delos #facebook #distributed #shared #log #algorithms #design #consensus #performance #scalability #faulttolerance #paper #research #infrastructure #storage #controlplane #corfu #paxos #zookeeper #zab

đź’» The byte is not enough
Fighting spam with Guardian, a real-time analytics and rules engine

Another great story of a practical approach to fight spam at big scale that is being applied at Pinterest. Though this isn't a story of inventing some tremendous and big all-in-one solution, it still provides a nice perspective on how to grow from a small and clumsy solution to a big and robust killer feature.

https://medium.com/pinterest-engineering/fighting-spam-with-guardian-a-real-time-analytics-and-rules-engine-938e7e61fa27

#systemdesign #infrastructure #pinterest #design #scalability #performance #bigdata #hive #kafka #presto #guardian #spam #realtime #analytics #safety #security

đź’» The byte is not enough
(Strong) Eventual Consistency and Conflict-free Replicated Data Types (CRDT)

Conflict-free Replicated Data Types (CRDTs) are an increasingly popular family of algorithms for optimistic replication. They allow data to be concurrently updated on several replicas, even while those replicas are offline, and provide a robust way of merging those updates back into a consistent state. CRDTs are used in geo-replicated databases, multi-user collaboration software, distributed processing frameworks, and various other systems.

However, while the basic principles of CRDTs are now quite well known, many challenging problems are lurking below the surface. It turns out that CRDTs are easy to implement badly. Many published algorithms have anomalies that cause them to behave strangely in some situations. Simple implementations often have terrible performance, and making the performance good is challenging.

Martin Kleppmann's talk on CRDT at Hydra Conf:
https://www.youtube.com/watch?v=PMVBuMK_pJY

Sean Cribbs discusses Convergent Replicated Data Types, data structures that tolerate eventual consistency:
https://www.infoq.com/presentations/CRDT/

Marc Shapiro presentation on SEC and CRTD at Microsoft Research:
https://www.youtube.com/watch?v=oyUHd894w18

#systemdesign #algorithms #sec #strong #consistency #crdt #replication #consensus #conflicts #resolution #riak #math #distributed

đź’» The byte is not enough
đź•° Blast from the past: Please stop calling databases CP or AP

A post from 2015 where Martin Kleppmann argues about usage of CAP-theorem in context of modern distributed systems. Martin does a great job demonstrating that people tend to attribute terms like consistency and availability with the meaning that differ from the one given in the original Brewer's CAP Theorem paper, and as thus demonstrates that almost no modern distributed system could be characterized with "A" or "C" property.

https://martin.kleppmann.com/2015/05/11/please-stop-calling-databases-cp-or-ap.html

#systemdesign #theory #research #cap #brewer #kleppmann #consistency #availability #partitiontolerance #zab #consensus #linearizability #zookeeper #riak #dynamodb #cassandra #voldemort #cp #aperture

đź’» The byte is not enough