Friday, May 21, 2021

System design: CAP theorem -> PACEL theorem | My 20 minutes study

May 21, 2021

My highlights:

  1. It states that in case of network partitioning (P) in a distributed computer system, one has to choose between availability (A) and consistency (C) (as per the CAP theorem), but else (E), even when the system is running normally in the absence of partitions, one has to choose between latency (L) and consistency (C)

PACELC theorem

From Wikipedia, the free encyclopedia
Jump to navigationJump to search

In theoretical computer science, the PACELC theorem is an extension to the CAP theorem. It states that in case of network partitioning (P) in a distributed computer system, one has to choose between availability (A) and consistency (C) (as per the CAP theorem), but else (E), even when the system is running normally in the absence of partitions, one has to choose between latency (L) and consistency (C).

Overview[edit]

PACELC builds on the CAP theorem. Both theorems describe how distributed databases have limitations and tradeoffs regarding consistency, availability, and partition tolerance. PACELC however goes further and states that an additional trade-off exists: between latency and consistency, even in absence of partitions, thus providing a more complete portrayal of the potential consistency trade-offs for distributed systems.[1]

A high availability requirement implies that the system must replicate data. As soon as a distributed system replicates data, a trade-off between consistency and latency arises.

The PACELC theorem was first described by Daniel J. Abadi from Yale University in 2010 in a blog post,[2] which he later clarified in a paper in 2012.[1] The purpose of PACELC is to address his thesis that "Ignoring the consistency/latency trade-off of replicated systems is a major oversight [in CAP], as it is present at all times during system operation, whereas CAP is only relevant in the arguably rare case of a network partition."

Database PACELC ratings[edit]

Database PACELC ratings are from [3]

  • The default versions of DynamoDB, Cassandra, Riak and Cosmos DB are PA/EL systems: if a partition occurs, they give up consistency for availability, and under normal operation they give up consistency for lower latency.
  • Fully ACID systems such as VoltDB/H-Store, Megastore and MySQL Cluster are PC/EC: they refuse to give up consistency, and will pay the availability and latency costs to achieve it. BigTable and related systems such as HBase are also PC/EC.
  • Couchbase provides a range of consistency and availability options during a partition, and equally a range of latency and consistency options with no partition. Unlike most other databases, Couchbase doesn't have a single API set nor does it scale/replicate all data services homogeneously. For writes, Couchbase favors Consistency over Availability making it formally CP, but on read there is more user-controlled variability depending on index replication, desired consistency level and type of access (single document lookup vs range scan vs full-text search, etc). On top of that, there is then further variability depending on cross-datacenter-replication (XDCR) which takes multiple CP clusters and connects them with asynchronous replication and Couchbase Lite which is an embedded database and creates a fully multi-master (with revision tracking) distributed topology.
  • Cosmos DB supports five tunable consistency levels that allow for tradeoffs between C/A during P, and L/C during E. Cosmos DB never violates the specified consistency level, so it’s formally CP.
  • MongoDB can be classified as a PA/EC system. In the baseline case, the system guarantees reads and writes to be consistent.
  • PNUTS is a PC/EL system.
  • Hazelcast IMDG and indeed most in-memory data grids are an implementation of a PA/EC system; Hazelcast can be configured to be EL rather than EC.[4] Concurrency primitives (Lock, AtomicReference, CountDownLatch, etc.) can be either PC/EC or PA/EC.[5]
  • FaunaDB implements Calvin, a transaction protocol created by Dr. Daniel Abadi and author[1] of PACELC theorem, and offers users adjustable controls for LC tradeoff. It is PC/EC for strictly serializable transactions, and EL for serializable reads.
DDBSP+AP+CE+LE+C
DynamoDBYesYes[a]
CassandraYesYes[a]
Cosmos DBYesYes
CouchbaseYesYesYes
RiakYesYes[a]
VoltDB/H-StoreYesYes
MegastoreYesYes
BigTable/HBaseYesYes
MySQL ClusterYesYes
MongoDBYesYes
PNUTSYesYes
Hazelcast IMDG[6][5]YesYesYesYes
FaunaDB[7]YesYesYes

See also[edit]

Notes[edit]

  1. ^ Jump up to:a b c Dynamo, Cassandra, and Riak have user-adjustable settings to control the LC tradeoff.[3]

References[edit]

  1. ^ Jump up to:a b c Abadi, Daniel J. "Consistency Tradeoffs in Modern Distributed Database System Design" (PDF). Yale University.
  2. ^ Abadi, Daniel J. (2010-04-23). "DBMS Musings: Problems with CAP, and Yahoo's little known NoSQL system". Retrieved 2016-09-11.
  3. ^ Jump up to:a b "Consistency Tradeoffs in Modern Distributed Database System Design" slide summary by Arinto Murdopo, Research Engineer
  4. ^ Abadi, Daniel (2017-10-08). "DBMS Musings: Hazelcast and the Mythical PA/EC System". DBMS Musings. Retrieved 2017-10-20.
  5. ^ Jump up to:a b "Hazelcast IMDG Reference Manual". docs.hazelcast.org. Retrieved 2020-09-17.
  6. ^ Abadi, Daniel (2017-10-08). "DBMS Musings: Hazelcast and the Mythical PA/EC System". DBMS Musings. Retrieved 2017-10-20.
  7. ^ Abadi, Daniel (2018-09-21). "DBMS Musings: NewSQL database systems are failing to guarantee consistency, and I blame Spanner". DBMS Musings. Retrieved 2019-02-23.

External links[edit]

HOW I GOT MY PHD | My Journey Through the Final (Crazy) Months!

 Here is the link.

NOK stock: NOK vs BB

May 21, 2021

Here is the article.

Nokia Corporation (NOK - Get Rating) is a Finland-based company engaged in the network and Internet protocol (IP) infrastructure, software, and related services market. The company’s networks segment comprises Mobile Access, Fixed Access, IP Routing, and Optical Networks businesses. NOK serves communications service providers, governments, large enterprises, and consumers.

BlackBerry Limited (BB - Get Rating) provides security software and services to enterprises and governments worldwide. The company leverages artificial intelligence (AI) and machine learning to deliver solutions in the areas of cybersecurity, safety and data privacy, and offers endpoint security management, encryption, and embedded systems.

The networking industry witnessed rising demand for its products and services over the past year owing to rapid adoption of remote working and learning. In fact, the global Data Center Networking market is expected to grow at a 7.2% CAGR over the next five years to reach $20.94 billion by 2025.

While NOK lost 2.2% over the past nine months, BB surged 83.1%. In terms of their past month’s performance, NOK is a clear winner with 19.1% gains versus BB’s 0.8% returns. But, which of these stocks is a better pick now? Let’s find out.

On May 18, 2021, NOK announced that it  will supply its Digital Operations software, cloud infrastructure software and AirFrame servers to help PLDT, a Philippines-based telecommunications company, and its wireless unit, Smart Communications, transform nationwide networks and standardize virtualization environments. With the help of NOK’s solutions, PLDT and Smart will be able to reduce operating costs and increase customer satisfaction by automating service, network and cloud operations. Its success should  result in a long-term partnership with NOK.

NOK was chosen by Telefónica’s Latin American brand, Movistar Chile, on May 10, to provide AirScale technology equipment for launching Movistar’s 5G network in Chile. NOK will also upgrade the company’s 4G and 4G+ networks to strengthen the critical network backbone across Chile’s key markets. Supported by Movistar’s fiber optics, NOK’s high-speed mobile technology is moving to deliver  next generation connectivity to enterprises and end users in Latin America.

On May 17, BB launched BlackBerry Optics 3.0, its next-generation cloud-based endpoint detection and response (EDR) solution, and BlackBerry Gateway, the company’s first AI-empowered Zero Trust Network Access (ZTNA) product. Available since the second quarter of 2021, the company hopes these products, rooted in a prevention-first and AI-driven approach, will provide enhanced visibility and protection against current and future cyberthreats.

Chinese electric carmaker WM Motor,  selected BB’s QNX software portfolio on May 6 to power its advanced W6 SUV model. BB QNX’s outstanding safety, cybersecurity and reliability enables WM Motors to focus on creating an extraordinary driving experience for its customers. This would make BB a trusted partner for the automotive industry.

Recent Financial Results

NOK’s net sales for its fiscal year 2021 first quarter, ended March 31, 2021, increased 3.3% year-over-year to €5.08 billion ($6.21 billion). Its net sales from the network infrastructure segment increased 21.7% year-over-year to €1.73 billion ($2.11 billion). The company’s gross profit came in at €1.93 billion ($2.36 billion), up 10.9% from the prior-year period. Its operating profit came in at €431 million ($578.73 million), compared to a loss of €76 million ($92.99 million) in the first quarter of 2020. Its net profit is reported €263 million ($321.79 million), compared to a loss of €115 million ($140.71 million) in the prior-year period. And its loss per share came in at €0.05, compared to a loss of €0.02 in the year-ago period.

For its fiscal year 2021 first quarter, ended February 28,  BB’s adjusted revenue declined 26.1% year-over-year to $215 million. The company’s adjusted gross profit has declined 29.1% year-over-year to $158 million. Its adjusted operating income came in at $18 million, which represents a 64.7% year-over-year decline. Its adjusted income came in at $16 million for the quarter, down 68.6% from the prior-year period. And, its adjusted EPS was $0.03, which represented a 66.7% year-over-year decline.

Past and Expected Financial Performance

NOK’s EBITDA grew at a 2,9% CAGR over the past three years. And the company’s total assets have declined at a rate of 2.1% over the past three years.

Analysts expect NOK’s revenue to increase 4.2% year-over-year for its fiscal year 2021 second quarter (ending June 30, 2021), marginally in the current year, and 2% in  2022. Its EPS is expected to decline by 30.9% year-over-year for the second quarter, but then increase 11% for the current year and 8.4% in 2022. NOK’s EPS is expected to grow at a rate of 16.5% per annum over the next five years.

In comparison, BB’s EBITDA grew at a 12.8% CAGR over the past three years. The company’s total assets have fallen at a rate of 9.3% over the past three years.

Analysts expect BB’s revenue to decline 20% in its  fiscal year 2021 second quarter (ending May 31, 2021), and 10.3% in 2022, but increase 15.4% in 2023. However, its EPS is expected to decrease 350% in the second quarter, 127.8% in the current year, and 220% next year. Furthermore, its EPS is expected to fall at a rate of 21.9% per annum over the next five years.

Profitability

NOK’s trailing-12-month revenue is 28.9 times  BB’s. NOK is also more profitable with a 10% gross profit margin versus RBLX’s negative value.

Also, NOK’s ROA and ROTC values of 3.6% and 6.5%, respectively, compare with BB’s negative values.

Valuation

In terms of forward non-GAAP P/E for the next fiscal year, BB is currently trading at 111.76x, 608.2% higher than NOK, which is currently trading at 15.78x. Also, NOK’s forward EV/sales of 0.94x is significantly lower than BB’s 6.15x. In terms of forward EV/EBITDA, BB’s 106.72x is 1328.6% higher than NOK’s 7.47x.

Thus, NOK looks more affordable here.

POWR Ratings

While BB has an overall D grade, which translates to Sell in our proprietary POWR Ratings system, NOK has an overall B grade, which equates to Buy. The POWR Ratings are calculated considering 118 different factors, each weighted to an optimal degree.

In terms of Sentiment, NOK has been graded an A, which is in sync with the company’s revenues and earnings growth potential expected by analysts. In comparison, BB’s Sentiment Grade of F is consistent with unfavorable analyst sentiment.

NOK has a B grade for Value. This is justified because its forward EV/EBITDA of 7.47x, which is 54.6% lower than the 16.45x industry average. However, BB has a C grade for Value. This is in sync with the company’s higher-than-industry forward EV/EBITDA value.

Of 55 stocks in the B-rated Technology – Communication/Networking industry, NOK is ranked #10 and BB is ranked #52.

Beyond what we’ve stated above, our POWR Ratings system has also rated both NOK and BB for Momentum, Stability, and Quality. Get all NOK ratings here. Also, click here to see the additional POWR Ratings for BB.

The Winner

The demand for networking is expected to remain high in 2021. As 5G technology is set to be available commercially this year, both NOK and BB are well-positioned to capitalize on the industry tailwinds. However, NOK appears to be a better buy based on its stable financials and higher profitability.

Our research shows that the odds of success increase if one bets on stocks with an Overall POWR Rating of Buy or Strong Buy. Click here to access the top-rated stocks in the Technology – Communication/Networking industry.

Thursday, May 20, 2021

System design: How distributed systems fail | written by Roberto Vitillo

 May 20, 2021

I plan to read the article and look into more detail if needed.


How distributed systems fail

December 05, 2020

At scale, any failure that can happen will eventually happen. Hardware failures, software crashes, memory leaks - you name it. The more components you have, the more failures you will experience.

This nasty behavior is caused by cruel math - given an operation that has a certain probability of failing, as the total number of operations performed increases, so does the total number of failures. In other words, as you scale out your application to handle more load, the more failures it will experience.

To protect your application against failures, you first need to know what can go wrong. Assuming you are using a cloud provider and not maintaning your own datacenter, the most common failures you will encounter are caused by single points of failure, the network being unreliable, slow processes, and unexpected load.

Single Point of Failure

A single point of failure is the most glaring cause of failure in a distributed system - it’s that one component that when it fails brings down the entire system with it. In practice, distributed systems can have multiple single points of failure.

A service that to start up needs to read its configuration from a non-replicated database is an example of a single point of failure - if the database isn’t reachable, the service won’t be able to start.

A more subtle example is a service that exposes a HTTP API on top of TLS and uses a certificate that needs to be manually renewed. If the certificate isn’t renewed by the time it expires, then most clients trying to connect to it wouldn’t be able to open a connection with the service.

Single points of failure should be identified when the system is architected before they can cause any harm. The best way to detect them is to examine every component of the system and ask what would happen if that component were to fail. Some single points of failure can be architected away, e.g., by introducing redundancy, while others can’t. In that case, the only option left is to minimize the blast radius.

Unreliable Network

When a client make a remote network call, it sends a request to a server and expects to receive a response from it a while later. In the best case, the client receives a response shortly after sending the request. But what if the client waits and waits and still doesn’t get a response?

In that case, the client doesn’t know whether a response will eventually arrive or not. At that point it has only two options, it can either continue to wait, or fail the request with an exception or an error.

Slow network calls are the silent killers of distributed systems. Because the client doesn’t know whether the response is on its way or not, it can spend a long time waiting before giving up, if it gives up at all. The wait can in turn cause degradations that are extremely hard to debug.

Slow Processes

From an observer’s point of view, a very slow process is not very different from one that isn’t running at all - neither can perform useful work. Resource leaks are one of the most common causes of slow processes.

Memory leaks are arguably the most well-known source of leaks. A memory leak manifests itself with a steady increase in memory consumption over time. Run-times with garbage collection don’t help much either - if a reference to an object that isn’t longer needed is kept somewhere, the object won’t be deleted by the garbage collector.

A memory leak keeps consuming memory until there is no more of it, at which point the operating system starts swapping memory pages to the disk constantly, all the while the garbage collector kicks in more frequently trying its best to release any shred of memory. The constant paging and the garbage collector eating up CPU cycles make the process slower. Eventually, when there is no more physical memory, and there is no more space in the swap file, the process won’t be able to allocate more memory, and most operations will fail.

Memory is just one of the many resources that can leak. For example, if you are using a thread pool, you can lose a thread when it blocks on a synchronous call that never returns. If a thread makes a synchronous, and blocking, HTTP call without setting a timeout, and the call never returns, the thread won’t be returned to the pool. Since the pool has a fixed size and keeps losing threads, the pool will eventually run out of threads.

You might think that making asynchronous calls, rather than a synchronous ones, would mitigate the problem in the previous case. But, modern HTTP clients use socket pools to avoid recreating TCP connections and pay a hefty performance fee. If a request is made without a timeout, the connection is never returned to the pool. As the pool has a limited size, eventually there won’t be any connections left to communicate with the host.

On top of all that, the code you write isn’t the only one accessing memory, threads and sockets. The libraries your application depends on access the same resources, and they can do all kinds of shady things. Without digging into their implementation, assuming it’s open in the first place, you can’t be sure whether they can wreak havoc or not.

Unexpected Load

Every system has a limit to how much load it can withstand without scaling. Depending on how the load increases, you are bound to hit that brick wall sooner or later. But one thing is an organic increase in load, which gives you the time to scale your service out accordingly, and another is a sudden and unexpected spike.

For example, consider the number of requests received by a service in a period of time. The rate and the type of incoming requests can change over time, and sometimes suddenly, for a variety of reasons:

  • The requests might have a seasonality - depending on the hour of the day the service is going to get hit by users in different countries.
  • Some requests are much more expensive than others and abuse the system in ways you didn’t really anticipate for, like scrapers slurping in data from your site at super human speed.
  • Some requests are malicious - think of DDoS attacks which try to saturate your service’s bandwidth, denying access to the service to legitimate users.

Cascading Failures

You would think that if your system has hundreds of processes, it shouldn’t make much of a difference if a small percentage are slow or unreachable. The thing about faults is that they tend to spread like cancer, propagating from one process to the other until the whole system crumbles to its knees. This effect is also referred to as a cascading failure, which occurs when a portion of an overall system fails, increasing the probability that other portions fail.

For example, suppose there are multiple clients querying two database replicas A and B, which are behind a load balancer. Each replica is handling about 50 transactions per second.

{width: 75%}cascading failure 1

Suddenly, replica B becomes unavailable because of a network fault. The load balancer detects that B is unavailable and removes it from its pool. Because of that, replica A has to pick up the slack for replica B, doubling the load it was previously under.

{width: 75%}cascading failure 2

As replica A starts to struggle to keep up with the incoming requests, the clients experience more failures and timeouts. In turn, they retry the same failing requests several times, adding insult to injury.

Eventually, replica A is under so much load that it can no longer serve requests promptly, and becomes for all intent and purposes unavailable, causing replica A to be removed from the load balancer’s pool. In the meantime, replica B becomes available again and the load balancer puts it back in the pool, at which point it’s flooded with requests that kill the replica instantaneously. This feedback loop of doom can repeat several time.

Cascading failures are very hard to get under control once they have started. The best way to mitigate one is to not have it in the first place by stopping the cracks in your services to propagate to others.

Defense Mechanisms

There is a variety of best practices you can use to mitigate failures, like circuit breakers, load shedding, rate-limiting and bulkheads. I plan to blog about those in the future, but in the meantime Google is your friend. Also, I have an entire chapter dedicated to resiliency patterns in my book about distributed systems.


Paper reading: Causal consistency | causal memory

May 20, 2021

  1. Traditional memory consistency? 
  2. causal memory, an abstraction that ensures that processes in a system agree on the relative ordering of operations that are causally related
  3. Because causal memory is weakly consistent, it admits more executions, and hence more concurrency, than either atomic or sequentially consistent memories

Causal Memory: 

Definitions, Implementation and Programming 

Mustaque Ahamad 

Gil Neiger 

James E. Burnsy 

Prince Kohli 

Georgia Institute of Technology 

Phillip W. Huttox 

GIT{CC{93/55 

September 17, 1993 

Revised: July 22, 1994 

Abstract 

The abstraction of a shared memory is of growing importance in distributed computing systems. Traditional memory consistency ensures that all processes agree on a common order of all operations on memory. Unfortunately, providing these guarantees entails access latencies that prevent scaling to large systems. This paper weakens such guarantees by defining causal memory, an abstraction that ensures that processes in a system agree on the relative ordering of operations that are causally related. Because causal memory is weakly consistent, it admits more executions, and hence more concurrency, than either atomic or sequentially consistent memories. This paper provides a formal definition of causal memory and gives an implementation for message-passing systems. In addition, it describes a practical class of programs that, if developed for a strongly consistent memory, run correctly with causal memory. 

1 Introduction 

The abstraction of a shared memory is of growing importance in distributed computing systems. It allows users to program these systems without concerning themselves with the details of the underlying message-passing system. Traditionally, distributed shared memories ensure that all processes in the system agree on a common order of all operations on memory. Such guarantees are provided by sequentially consistent memory [27] and by atomic memory [28] (also called linearizable memory [20]). Unfortunately, providing these consistency guarantees entails access latencies that prevent scaling to large systems. A simple argument [10,29] can be used to show that no memory can provide strong consistency and retain low latency in systems with high message-passing delays. This tradeoff represents a significant effciency problem since it forces applications to pay the costs of consistency even if they are highly parallel and involve little synchronization. A number of techniques [11,24] have been suggested to improve the effciency of shared memory implementations, but all provide only partial remedies to the fundamental problem of latency and scale for strongly consistent memories. 

Recent research [1,6,9,16{18,21,29] suggests that a systematic weakening of memory consistency can reduce the costs of providing consistency while maintaining a viable target" model for programmers. Weakly consistent memories admit more executions, and hence more concurrency, than either sequentially consistent or atomic memories. This paper defines causal memory, an abstraction that ensures that processes in a system agree on the relative ordering of operations that are causally related. (Causal memory has been mentioned earlier [6,21]; however, these papers do not present careful definitions as is done here.) This paper provides a formal definition of causal memory and gives an implementation for message-passing systems. We give two classes of programs that can be developed assuming a sequentially consistent memory and that run correctly with causal memory.

Causal memory is based on Lamport's concept of potential causality [26]. Potential causality provides a natural ordering on events in a distributed system where processes communicate via message passing. We introduce a similar notion of causality based on reads and writes in a shared memory environment. Causal memory requires that reads return values consistent with causally related reads and writes, and we say that "reads respect the order of causally related writes." Since causality orders events only partially, reading processes may disagree on the relative ordering of concurrent writes. This provides independence between concurrent writers which reduces consistency maintenance (synchronization) costs. The idea is that the synchronization required by a program is often specified explicitly and it is not necessary for the memory to provide additional synchronization guarantees.

Causal memory is related to the ISIS causal broadcast and, thereby, to the notion of causally ordered messages [13]. Our implementation of causal memory is based on the use of vector timestamps [14,30], as is the ISIS implementation of causal broadcast. Both implementations are "non-blocking": a process may complete an operation (e.g., a write or a send) without waiting for communication with other processes. Nevertheless, causal memory is more than a collection of "locations" updated by causal broadcasts. Memory has overwrite semantics and messages have queuing semantics. A message recipient can be assured that it will eventually receive all messages that have been sent to it, but repeated reads cannot guarantee that all values written will be read. "Hidden writes", values overwritten before they are read, are always possible. Since a process may read memory locations in any order it chooses, it may read a value v1 from location x much later than a value v2 from location y, even when the write operation that stores v1 in x is causally before the write of v2 to y. In a message-passing system, such behavior would violate the required causal ordering.

We give precise characterizations of two classes of programs that run correctly with causal memory. Any execution a program in either of these classes with causal memory is actually sequentially consistent. If the program is proven correct with sequential consistent memory, then it is still correct with causal memory. One of these classes includes data-race free programs [1,2] that make use of explicit synchronization to prevent problems that may stem from concurrent access to shared memory. 

It is far from clear that there is a "best" kind of shared memory model for use with distributed systems. Strongly consistent memories are easier to program than weak memories, but they require costly blocking implementations. Very weak memories may be implemented cheaply, but they might not be practical to program. We believe that causal memory provides a happy medium: it allows non-blocking implementations and is a useful model for a class of practical programs.



System design: Distributed database consistency model | Causal consistency | Happen before | 用happens before关系的有向无环图定义了Causal consistency

 

深入浅析一致性模型之Causal Consistency

2019.05.25 19:16 3302浏览

本文是《如何学习分布式系统》中,关于一致性模型的相关介绍。

Causal consistency的定义

Causal consistency叫做因果一致性,被认为是比Sequential Consistency更弱的一致性,因为在Causal consistency中,只对有因果关系的事件有顺序要求。

Causal consistency的概念最早在《Causal memory: Definitions, implementation, and programming》一文中被提出。随后在《Consistency, Availability, and Convergence》一文中,作者用happens before关系的有向无环图定义了Causal consistency,并且提出了real time causal consistency一致性模型。这篇文章得出了以下两个重要的结论,有兴趣的同学可以读一下原文:

  1. No consistency stronger than real time causal (RTC) consistency, a strengthening of causal consistency, can be provided in an always-available, one-way convergent system.
  2. RTC can be provided in an always-available, one-way convergent system.

Causal consistency要求如果两个事件有因果关系,那么在所有节点上必须观测到这个因果关系。

Causal consistency的例子

比如下图中,我们认为P2写入的3是基于它读出来的1计算出来的,它读出来的1又是由P1的写入产生的,因此认为P1写入1和P2写入3具有因果关系。P4没有观测到这个因果关系,所以这个系统不具备Causal Consistency。

Causal Consistency例子1

而下图中,认为P2写入3和P1写入1不具有因果关系,则P4和P3可以以任意顺序观测到它们。这个系统仍然可以说具有Causal consistency,但是不具备Sequential Consistency。
Causal Consistency例子2

Causal consistency的应用

Causal consistency一般应用在跨地域同步数据中心系统中,例如Facebook、微信这样的应用程序,全球各地的用户,往往会访问其距离最近的数据中心,数据中心之间再进行双向的数据同步。为了减小数据同步的延迟,往往并行的同步数据。
图片描述
没有因果一致性时会发生如下情形:

  1. 夏侯铁柱在朋友圈发表状态“我戒指丢了”
  2. 夏侯铁柱在同一条状态下评论“我找到啦”
  3. 诸葛建国在同一条状态下评论“太棒了”
  4. 远在美国的键盘侠看到“我戒指丢了”“太棒了”,开始喷诸葛建国
  5. 远在美国的键盘侠看到“我戒指丢了”“我找到啦”“太棒了”,意识到喷错人了

或者:

  1. 夏侯铁柱从好友中删除了诸葛建国
  2. 夏侯铁柱发表了朋友圈“清理了一波没用的好友”
  3. 远在美国的诸葛建国看到了该朋友圈
  4. 诸葛建国想去点个赞。。。

所以很多系统采用因果一致性系统来避免这种问题,我们将会在后续的文章中介绍。

例如微信的朋友圈就采用了因果一致性,但是它的资料有点过于简略,有兴趣的可以参考https://www.useit.com.cn/thread-10587-1-1.html

Walk 2 KM | Soaked in the rain | Healthy life style | May 17, 2021