Learn Labs
13. A Philosophy of Streaming Systems

13.5 Aiming for Correctness

"With STATELESS services that only read data, it's NOT A BIG DEAL if something goes wrong; you can fix the bug and restart. STATEFUL systems are not so simple. THEY ARE DESIGNED TO REMEMBER THINGS FOREVER, SO IF SOMETHING GOES WRONG, THE EFFECTS ALSO POTENTIALLY LAST FOREVER."

The unflattering state of the art:

"For approximately four decades, ATOMICITY, ISOLATION, AND DURABILITY have been the tools of choice. HOWEVER, THOSE FOUNDATIONS ARE WEAKER THAN THEY SEEM: witness the confusion of weak isolation levels."

"CONSISTENCY IS OFTEN TALKED ABOUT BUT POORLY DEFINED. Some people assert that we should 'EMBRACE WEAK CONSISTENCY' for the sake of better availability, WHILE LACKING A CLEAR IDEA OF WHAT THAT MEANS IN PRACTICE."

"For a topic that is so important, OUR UNDERSTANDING AND OUR ENGINEERING METHODS ARE SURPRISINGLY FLAKY. It is VERY DIFFICULT TO DETERMINE WHETHER IT IS SAFE to run a particular application using a particular isolation level or replication configuration. Often, SIMPLE SOLUTIONS APPEAR TO WORK CORRECTLY WHEN CONCURRENCY IS LOW AND THERE ARE NO FAULTS, BUT TURN OUT TO HAVE MANY SUBTLE BUGS IN MORE DEMANDING CIRCUMSTANCES."

"Kyle Kingsbury's JEPSEN experiments have highlighted THE STARK DISCREPANCIES BETWEEN SOME PRODUCTS' CLAIMED SAFETY GUARANTEES AND THEIR ACTUAL BEHAVIOR. Even if infrastructure products were free from problems, APPLICATION CODE WOULD STILL NEED TO CORRECTLY USE THE FEATURES THEY PROVIDE, WHICH IS ERROR-PRONE IF THE CONFIGURATION IS HARD TO UNDERSTAND."

5.1 The end-to-end argument

"JUST BECAUSE AN APPLICATION USES A DATA SYSTEM THAT PROVIDES COMPARATIVELY STRONG SAFETY PROPERTIES, SUCH AS SERIALIZABLE TRANSACTIONS, THAT DOES NOT MEAN THE APPLICATION IS GUARANTEED TO BE FREE FROM DATA LOSS OR CORRUPTION. If an application has a bug that causes it to write incorrect data or delete data, SERIALIZABLE TRANSACTIONS AREN'T GOING TO SAVE YOU."

"This is an argument in favor of IMMUTABLE AND APPEND-ONLY DATA, because it is easier to recover from such mistakes if YOU REMOVE THE ABILITY OF FAULTY CODE TO DESTROY GOOD DATA."

The four-layer duplicate-suppression failure — trace it carefully:

  1. TCP. TCP uses sequence numbers to put packets in order and determine whether any were lost or duplicated. Lost packets are retransmitted and duplicates removed before it hands the data to an application. ✗ But this works only within the context of a single TCP connection.
  2. A database transaction.
    BEGIN TRANSACTION;
    UPDATE accounts SET balance = balance + 11.00 WHERE account_id = 1234;
    UPDATE accounts SET balance = balance - 11.00 WHERE account_id = 4321;
    COMMIT;
    ✗ If the client suffers a network interruption after sending the commit but before hearing back, it does not know whether the transaction committed or aborted. The client can reconnect and retry — but now it is outside the scope of TCP duplicate suppression. Since this transaction is not idempotent, $22 could be transferred instead of the desired $11. Thus, even though code like this is a standard example for transaction atomicity, it is not correct, and real banks do not work like this.
  3. Two-phase commit. 2PC breaks the one-to-one mapping between a TCP connection and a transaction, since the coordinator must be able to reconnect and tell the database whether to commit or abort. Is this sufficient? ✗ Unfortunately not.
  4. The end-user device — where it actually breaks. The user has a weak cellular connection. They succeed in sending the POST request but lose the signal before receiving the response. They will be shown an error, and they may retry manually. Web browsers warn, “are you sure you want to submit this form again?” — and the user says yes, because they want the operation to happen. From the web server's point of view the retry is a separate request, and from the database's point of view a separate transaction. The usual deduplication mechanisms don't help.

The fix — a request ID passed end to end:

ALTER TABLE requests ADD UNIQUE (request_id);

BEGIN TRANSACTION;
INSERT INTO requests
  (request_id, from_account, to_account, amount)
  VALUES('0286FDB8-D7E1-423F-B40B-792B3608036C', 4321, 1234, 11.00);
UPDATE accounts SET balance = balance + 11.00 WHERE account_id = 1234;
UPDATE accounts SET balance = balance - 11.00 WHERE account_id = 4321;
COMMIT;

"Generate a unique identifier for each request (such as a UUID) and include it AS A HIDDEN FORM FIELD in the client application, or CALCULATE A HASH OF ALL THE RELEVANT FORM FIELDS to derive the request ID. IF THE BROWSER SUBMITS TWICE, THE TWO REQUESTS WILL HAVE THE SAME REQUEST ID."

"This relies on A UNIQUENESS CONSTRAINT. RELATIONAL DATABASES CAN GENERALLY MAINTAIN A UNIQUENESS CONSTRAINT CORRECTLY, EVEN AT WEAK ISOLATION LEVELS — whereas an APPLICATION-LEVEL CHECK-THEN-INSERT MAY FAIL under nonserializable isolation."

Bonus: "the requests table acts as A KIND OF EVENT LOG, useful for event sourcing or CDC. The updates to the balances DON'T HAVE TO HAPPEN IN THE SAME TRANSACTION, since they are REDUNDANT AND COULD BE DERIVED FROM THE REQUEST EVENT in a downstream consumer — as long as the event is processed exactly once, WHICH CAN AGAIN BE ENFORCED USING THE REQUEST ID."

The principle, stated in 1984 by Saltzer, Reed, and Clark:

"THE FUNCTION IN QUESTION CAN COMPLETELY AND CORRECTLY BE IMPLEMENTED ONLY WITH THE KNOWLEDGE AND HELP OF THE APPLICATION STANDING AT THE ENDPOINTS OF THE COMMUNICATION SYSTEM. THEREFORE, PROVIDING THAT QUESTIONED FUNCTION AS A FEATURE OF THE COMMUNICATION SYSTEM ITSELF IS NOT POSSIBLE. (Sometimes an incomplete version provided by the communication system MAY BE USEFUL AS A PERFORMANCE ENHANCEMENT.)"

It generalizes to three things:

FunctionLow-level mechanismWhy it's insufficient
Duplicate suppressionTCP sequence numbers; stream processor exactly-onceCan't prevent a user resubmitting a timed-out form
Integrity checkingChecksums in Ethernet, TCP, TLS"They cannot detect corruption due to BUGS IN THE SOFTWARE at the sending and receiving ends, or CORRUPTION ON THE DISKS." ⇒ you need END-TO-END CHECKSUMS
EncryptionWiFi password protects against local snooping; TLS protects against network attackers"Neither protects against COMPROMISES OF THE SERVER. Only END-TO-END ENCRYPTION AND AUTHENTICATION can protect against all these things"

"Although the low-level features CANNOT PROVIDE THE DESIRED END-TO-END FEATURES BY THEMSELVES, THEY ARE STILL USEFUL, SINCE THEY REDUCE THE PROBABILITY OF PROBLEMS AT HIGHER LEVELS. HTTP requests would often get mangled if we didn't have TCP putting packets back in order. WE JUST NEED TO REMEMBER THAT THE LOW-LEVEL RELIABILITY FEATURES ARE NOT BY THEMSELVES SUFFICIENT."

The uncomfortable conclusion:

"That is a shame, because FAULT-TOLERANCE MECHANISMS ARE HARD TO GET RIGHT. It would be really nice to WRAP UP THE HIGH-LEVEL FAULT-TOLERANCE MACHINERY IN AN ABSTRACTION so that application code needn't worry about it — BUT IT SEEMS THAT WE HAVE NOT YET FOUND THE RIGHT ONE."

"Transactions collapse a wide range of issues down to TWO POSSIBLE OUTCOMES: COMMIT OR ABORT. That is A HUGE SIMPLIFICATION — BUT IT IS NOT ENOUGH. Transactions are EXPENSIVE, especially with heterogeneous storage. WHEN WE REFUSE TO USE DISTRIBUTED TRANSACTIONS BECAUSE THEY ARE TOO EXPENSIVE, WE END UP HAVING TO REIMPLEMENT FAULT-TOLERANCE MECHANISMS IN APPLICATION CODE. As numerous examples have shown, REASONING ABOUT CONCURRENCY AND PARTIAL FAILURE IS DIFFICULT AND COUNTERINTUITIVE, AND SO MOST APPLICATION-LEVEL MECHANISMS DO NOT WORK CORRECTLY. THE CONSEQUENCE IS LOST OR CORRUPTED DATA."

5.2 Enforcing constraints without distributed transactions

Uniqueness requires consensus — and here's how to get it from a log:

① append the request②taken?emit③ client watchesclientlog shardhash of the usernamestream processorreads sequentiallylocal databasewhich usernames are takenoutput streamsuccess or rejectionavailable → record it as taken, emit a success message · taken → emit a rejection message

This algorithm is the same as the construction for achieving consensus using a shared log (Ch 10). It scales easily to a large request throughput by increasing the number of shards, as each shard can be processed independently.

The general principle: any writes that may conflict are routed to the same shard and processed sequentially. The definition of a conflict may depend on the application, but the stream processor can use arbitrary logic to validate.

Asynchronous multi-leader replication is ruled out, because different leaders could concurrently accept conflicting writes. If you want to immediately reject any writes that would violate the constraint, synchronous coordination is unavoidable.

Figure 13.5.2Uniqueness requires consensus — and here's how to get it from a log

Multishard atomicity WITHOUT atomic commit — the money-transfer worked example:

The traditional way is an atomic commit across three shards — request ID, payer, payee — which essentially forces it into a total order with respect to all other transactions on any of those shards. Since there is now cross-shard coordination, different shards can no longer be processed independently, so throughput is likely to suffer.

The dataflow way instead has the source-account processor maintain a local database with the account state and the IDs of requests it has already processed — entirely derived from the log. On an unseen ID it checks whether there is enough money, and if so reserves the amount locally and emits three events, all carrying the original request ID.

clientsource shardsource processordestination shardfees shard① transfer request + request ID② read in log orderoutgoing paymentincoming paymentincoming payment③ delivered back ⇒ execute

The outgoing payment event goes to the source shard, which is the processor's own input log, so it is eventually delivered back. The processor recognizes, based on the request ID, that this is a payment it previously reserved and executes it, ignoring duplicates on the same ID. ④ The destination and fees shards are consumed by independent tasks, which update their local state and deduplicate on request ID.

Crash analysis: if the source processor crashes mid-request, the output messages may or may not have been emitted. After recovery it processes the same request again (at-least-once) and makes the same decision (deterministic). It emits the same output messages with the same request ID. If they are duplicates, downstream consumers ignore them.

Atomicity in this system comes not from transactions, but from the fact that writing the initial request event to the source account log is an atomic action. Once that one event is in the log, all the downstream events will eventually be written as well — possibly after recovery, possibly with duplicates, but they will appear eventually.

The only requirements are that events for any given account are processed strictly in log order with at-least-once semantics, and that the stream processors are deterministic. The three accounts could just as well be in the same shard — it doesn't matter.

Figure 13.5.3Multishard atomicity WITHOUT atomic commit — the money-transfer worked example

5.3 Timeliness vs Integrity — the chapter's most important distinction

TimelinessIntegrity
Ensuring that users observe the system in an up-to-date state.Absence of corruption — no data loss, and no contradictory or false data. If a derived dataset is a view onto underlying data, the derivation must be correct.
That inconsistency is temporary, and it will eventually be resolved simply by waiting and trying again.If integrity is violated, the inconsistency is permanent; waiting and trying again is not going to fix database corruption. Instead, explicit checking and repair is needed.
CAP's “consistency” is linearizability, a strong way of achieving timeliness. Weaker forms include read-after-write consistency.ACID's “consistency” is an application-specific notion of integrity. Atomicity and durability are important tools for preserving integrity.

In slogan form: violations of timeliness are allowed under eventual consistency, whereas violations of integrity result in perpetual inconsistency.

In most applications, integrity is much more important than timeliness. Violations of timeliness can be annoying and confusing, but violations of integrity can be catastrophic.

The credit card statement example, which makes it concrete:

"It is NOT SURPRISING if a transaction you made within the last 24 hours DOES NOT YET APPEAR. It is normal that these systems have a certain lag. WE KNOW THAT BANKS RECONCILE AND SETTLE TRANSACTIONS ASYNCHRONOUSLY, AND TIMELINESS IS NOT VERY IMPORTANT HERE."

"HOWEVER, IT WOULD BE VERY BAD IF THE STATEMENT BALANCE WAS NOT EQUAL TO THE SUM OF THE TRANSACTIONS PLUS THE PREVIOUS BALANCE (an error in the sums), OR IF A TRANSACTION WAS CHARGED TO YOU BUT NOT PAID TO THE MERCHANT (DISAPPEARING MONEY). SUCH PROBLEMS WOULD BE VIOLATIONS OF THE INTEGRITY OF THE SYSTEM."

Why this matters so much for dataflow:

"ACID transactions usually provide BOTH timeliness AND integrity. Thus, if you approach correctness from the point of view of ACID, THE DISTINCTION IS FAIRLY INCONSEQUENTIAL."

"AN INTERESTING PROPERTY OF EVENT-BASED DATAFLOW SYSTEMS IS THAT THEY DECOUPLE TIMELINESS AND INTEGRITY. When processing asynchronously, THERE IS NO GUARANTEE OF TIMELINESS unless you explicitly build consumers that wait. HOWEVER, INTEGRITY IS IN FACT CENTRAL TO STREAMING SYSTEMS. Exactly-once semantics IS A MECHANISM FOR PRESERVING INTEGRITY."

The four mechanisms that give integrity without distributed transactions:

  1. "REPRESENTING THE CONTENT OF THE WRITE OPERATION AS A SINGLE MESSAGE, which can easily be written atomically — an approach that FITS VERY WELL WITH EVENT SOURCING"
  2. "DERIVING ALL OTHER STATE UPDATES FROM THAT SINGLE MESSAGE VIA DETERMINISTIC DERIVATION FUNCTIONS, similarly to stored procedures"
  3. "PASSING A CLIENT-GENERATED REQUEST ID THROUGH ALL THESE LEVELS OF PROCESSING, enabling END-TO-END DUPLICATE SUPPRESSION AND IDEMPOTENCE"
  4. "MAKING MESSAGES IMMUTABLE AND ALLOWING DERIVED DATA TO BE REPROCESSED from time to time, WHICH MAKES IT EASIER TO RECOVER FROM BUGS"

5.4 Loosely interpreted constraints — the business-reality argument

"Enforcing a uniqueness constraint requires consensus, typically implemented by funneling all events through a single node. THIS LIMITATION IS UNAVOIDABLE if we want the traditional form. HOWEVER, MANY REAL APPLICATIONS HAVE A BUSINESS REQUIREMENT TO ALLOW VIOLATIONS OF WHAT YOU MIGHT THINK OF AS HARD CONSTRAINTS:"

CaseThe reality
Overselling stock"You can order in more stock, apologize for the delay, and offer a discount. THIS IS THE SAME AS WHAT YOU'D HAVE TO DO IF A FORKLIFT TRUCK RAN OVER SOME OF THE ITEMS IN YOUR WAREHOUSE. THUS, THE APOLOGY WORKFLOW ALREADY NEEDS TO BE PART OF YOUR BUSINESS PROCESSES ANYWAY, AND A HARD CONSTRAINT MIGHT BE UNNECESSARY"
Overbooking flights and hotels"The constraint of 'one person per seat' is DELIBERATELY VIOLATED FOR BUSINESS REASONS, and compensation processes (refunds, upgrades, a complimentary room at a neighboring hotel) are put in place. EVEN IF NO OVERBOOKING OCCURRED, APOLOGY AND COMPENSATION PROCESSES WOULD BE NEEDED to deal with flights canceled because of BAD WEATHER OR STAFF GOING ON STRIKE"
Overdrafts"The bank can charge an overdraft fee and ask them to pay back what they owe. BY LIMITING TOTAL WITHDRAWALS PER DAY, THE RISK TO THE BANK IS BOUNDED"
Cross-organization integration"Inconsistencies WILL INEVITABLY ARISE, and correction mechanisms are necessary. Settlement of payments between banks is an example"

A change to correct a mistake is a COMPENSATING TRANSACTION. "The cost of the apology varies, but IT IS OFTEN QUITE LOW; YOU CAN'T UNSEND AN EMAIL, BUT YOU CAN SEND A FOLLOW-UP EMAIL WITH A CORRECTION. If you accidentally charge a credit card twice, you can refund one, and the cost is just the processing fees and perhaps a customer complaint. Once money has been paid out of an ATM you can't directly get it back — although in principle you can SEND DEBT COLLECTORS."

"IF THE COST OF THE APOLOGY IS ACCEPTABLE, THE TRADITIONAL MODEL OF CHECKING ALL CONSTRAINTS BEFORE EVEN WRITING THE DATA IS UNNECESSARILY RESTRICTIVE. It may well be reasonable to GO AHEAD WITH A WRITE OPTIMISTICALLY AND CHECK THE CONSTRAINT AFTER THE FACT. You can still ensure that VALIDATION OCCURS BEFORE TAKING ACTIONS THAT WOULD BE EXPENSIVE TO RECOVER FROM, BUT THAT DOESN'T IMPLY YOU MUST VALIDATE BEFORE YOU EVEN WRITE THE DATA."

"These applications DO require integrity. You would not want to lose a reservation or have money disappear. BUT THEY DON'T REQUIRE TIMELINESS ON THE ENFORCEMENT OF THE CONSTRAINT."

5.5 Coordination-avoiding data systems

The two observations, combined:

  1. “Dataflow systems can maintain Integrity guarantees on derived data Without atomic commit, linearizability, or synchronous cross-shard coordination.”
  2. “Although strict uniqueness constraints require timeliness and coordination, Many applications are fine with loose constraints that may be temporarily violated and fixed up later, As long as integrity is preserved throughout.”

⇒ “Coordination-avoiding data systems can provide data management services for many applications Without requiring coordination, while still giving strong integrity guarantees. They can achieve Better performance and fault tolerance than systems that need synchronous coordination.”

“For example, such a system could operate across multiple datacenters in a Multi-leader configuration, asynchronously replicating between regions. Any one datacenter can continue operating independently. Such a system would have Weak timeliness guarantees — it could not be linearizable without introducing coordination — but it can still have strong integrity guarantees.”

“Serializable transactions are Still useful as part of maintaining derived state, but They can be run at a small scope where they work well. Heterogeneous distributed transactions such as XA Are not required. Synchronous coordination can still be introduced Where it is needed (to enforce strict constraints before an operation from which recovery is not possible), But there is no need for everything to pay the cost of coordination if only a small part of an application needs it.”

THE APOLOGY CALCULUS: "Coordination and constraints REDUCE THE NUMBER OF APOLOGIES YOU HAVE TO MAKE FOR INCONSISTENCIES, BUT POTENTIALLY ALSO REDUCE THE PERFORMANCE AND AVAILABILITY OF YOUR SYSTEM, AND THUS POTENTIALLY INCREASE THE NUMBER OF APOLOGIES YOU HAVE TO MAKE FOR OUTAGES. YOU CANNOT REDUCE THE NUMBER OF APOLOGIES TO ZERO, BUT YOU CAN AIM TO FIND THE BEST TRADE-OFF — THE SWEET SPOT WITH NEITHER TOO MANY INCONSISTENCIES NOR TOO MANY AVAILABILITY PROBLEMS."

5.6 Trust, but Verify

"Traditionally, system models take A BINARY APPROACH toward faults: we assume that some things can happen and that other things CAN NEVER happen. IN REALITY, IT IS MORE A QUESTION OF PROBABILITIES. The question is WHETHER VIOLATIONS OF OUR ASSUMPTIONS HAPPEN OFTEN ENOUGH THAT WE MAY ENCOUNTER THEM IN PRACTICE."

"We have seen that data can become corrupted IN MEMORY, ON DISK, AND ON THE NETWORK. MAYBE THIS IS SOMETHING WE SHOULD BE PAYING MORE ATTENTION TO? IF YOU ARE OPERATING AT LARGE ENOUGH SCALE, EVEN VERY UNLIKELY THINGS DO HAPPEN."

Software bugs in the systems you trust most:

"Even widely used database software has bugs — PAST VERSIONS OF MySQL HAVE FAILED TO CORRECTLY MAINTAIN UNIQUENESS CONSTRAINTS, AND POSTGRESQL'S SERIALIZABLE ISOLATION LEVEL HAS EXHIBITED WRITE SKEW ANOMALIES IN THE PAST — even though MySQL and PostgreSQL are ROBUST AND WELL-REGARDED databases BATTLE-TESTED BY MANY PEOPLE FOR MANY YEARS. IN LESS MATURE SOFTWARE, THE SITUATION IS LIKELY TO BE MUCH WORSE."

"When it comes to APPLICATION code, we have to assume MANY MORE BUGS, since most applications don't receive anywhere near the amount of review and testing that database code does. MANY APPLICATIONS DON'T EVEN CORRECTLY USE THE FEATURES DATABASES OFFER FOR PRESERVING INTEGRITY, SUCH AS FOREIGN-KEY OR UNIQUENESS CONSTRAINTS."

"ACID consistency is based on the idea that the database starts in a consistent state and a transaction transforms it to another consistent state. HOWEVER, THIS NOTION MAKES SENSE ONLY IF WE ASSUME THE TRANSACTION IS FREE FROM BUGS. IF THE APPLICATION USES THE DATABASE INCORRECTLY — FOR EXAMPLE, USING A WEAK ISOLATION LEVEL UNSAFELY — THE INTEGRITY OF THE DATABASE CANNOT BE GUARANTEED."

What mature systems actually do:

"Large-scale storage systems such as HDFS AND AMAZON S3 DO NOT FULLY TRUST DISKS. These systems run BACKGROUND PROCESSES THAT CONTINUALLY READ BACK FILES, COMPARE THEM TO OTHER REPLICAS, AND MOVE FILES FROM ONE DISK TO ANOTHER, in order to mitigate the risk of SILENT CORRUPTION."

"IF YOU WANT TO BE SURE THAT YOUR DATA IS STILL THERE, YOU HAVE TO READ IT AND CHECK. By the same argument, IT IS IMPORTANT TO TRY RESTORING FROM YOUR BACKUPS FROM TIME TO TIME — OTHERWISE YOU MAY FIND OUT THAT YOUR BACKUP IS BROKEN WHEN IT IS TOO LATE AND YOU HAVE ALREADY LOST DATA. DON'T JUST BLINDLY TRUST THAT IT IS ALL WORKING."

"NOT MANY SYSTEMS CURRENTLY HAVE THIS KIND OF 'TRUST, BUT VERIFY' APPROACH OF CONTINUALLY AUDITING THEMSELVES. MANY ASSUME THAT CORRECTNESS GUARANTEES ARE ABSOLUTE AND MAKE NO PROVISION FOR THE POSSIBILITY OF RARE DATA CORRUPTION. In the future we may see more SELF-VALIDATING OR SELF-AUDITING SYSTEMS."

Designing for auditability:

"If a transaction mutates several objects, THE UNDERLYING REASON CAN BE DIFFICULT TO TELL AFTER THE FACT. Even if you capture the transaction logs, THE INSERTIONS, UPDATES, AND DELETIONS DO NOT NECESSARILY GIVE A CLEAR PICTURE OF WHY those mutations were performed. THE INVOCATION OF THE APPLICATION LOGIC THAT DECIDED ON THOSE MUTATIONS IS TRANSIENT AND CANNOT BE REPRODUCED."

"By contrast, EVENT-BASED SYSTEMS CAN PROVIDE BETTER AUDITABILITY. User input is represented as a SINGLE IMMUTABLE EVENT, and any resulting state updates are DERIVED from it. The derivation can be made DETERMINISTIC AND REPEATABLE."

"Being explicit about dataflow MAKES THE PROVENANCE OF DATA MUCH CLEARER. For the event log, we can USE HASHES to check that the event storage has not been corrupted. For any derived state, we can RERUN THE BATCH AND STREAM PROCESSORS to check whether we get the same result, OR EVEN RUN A REDUNDANT DERIVATION IN PARALLEL."

"A deterministic and well-defined dataflow also makes it easier to DEBUG AND TRACE the execution to determine WHY it did something — A KIND OF TIME-TRAVEL DEBUGGING CAPABILITY."

And the end-to-end argument, one more time:

"CHECKING THE INTEGRITY OF DATA SYSTEMS IS BEST DONE IN AN END-TO-END FASHION. THE MORE SYSTEMS WE CAN INCLUDE IN AN INTEGRITY CHECK, THE FEWER OPPORTUNITIES THERE ARE FOR CORRUPTION TO GO UNNOTICED. IF WE CAN CHECK THAT AN ENTIRE DERIVED DATA PIPELINE IS CORRECT END TO END, THEN ANY DISKS, NETWORKS, SERVICES, AND ALGORITHMS ALONG THE PATH ARE IMPLICITLY INCLUDED IN THE CHECK."

"Having continuous end-to-end integrity checks GIVES YOU INCREASED CONFIDENCE ABOUT THE CORRECTNESS OF YOUR SYSTEMS, WHICH IN TURN ALLOWS YOU TO MOVE FASTER. LIKE AUTOMATED TESTING, AUDITING INCREASES THE CHANCES THAT BUGS WILL BE FOUND QUICKLY. IF YOU ARE NOT AFRAID OF MAKING CHANGES, YOU CAN MUCH BETTER EVOLVE AN APPLICATION."

Tools — and the blockchain connection, treated soberly:

"A transaction log can be made TAMPER-PROOF by periodically signing it with a hardware security module, BUT THAT DOES NOT GUARANTEE THAT THE RIGHT TRANSACTIONS WENT INTO THE LOG IN THE FIRST PLACE."

"BLOCKCHAINS are SHARED APPEND-ONLY LOGS WITH CRYPTOGRAPHIC CONSISTENCY CHECKS; THE TRANSACTIONS THEY STORE ARE EVENTS, AND SMART CONTRACTS ARE BASICALLY STREAM PROCESSORS. The difference from the consensus protocols of Ch 10 is that blockchains are BYZANTINE FAULT-TOLERANT — they still work if some nodes have corrupted data BECAUSE THE REPLICAS CONTINUALLY CHECK ONE ANOTHER'S INTEGRITY."

"FOR MOST APPLICATIONS, BLOCKCHAINS HAVE TOO HIGH AN OVERHEAD TO BE USEFUL. HOWEVER, SOME OF THEIR CRYPTOGRAPHIC TOOLS CAN BE USED IN A LIGHTER-WEIGHT CONTEXT. MERKLE TREES are trees of hashes that can efficiently prove a record appears in a dataset. CERTIFICATE TRANSPARENCY uses cryptographically verified append-only logs and Merkle trees to check TLS/SSL certificates; IT AVOIDS NEEDING A CONSENSUS PROTOCOL BY HAVING A SINGLE LEADER PER LOG."


On this page

"For a topic that is so important, OUR UNDERSTANDING AND OUR ENGINEERING METHODS ARE SURPRISINGLY FLAKY. It is VERY DIFFICULT TO DETERMINE WHETHER IT IS SAFE to run a particular application using a particular isolation level or replication configuration. Often, SIMPLE SOLUTIONS APPEAR TO WORK CORRECTLY WHEN CONCURRENCY IS LOW AND THERE ARE NO FAULTS, BUT TURN OUT TO HAVE MANY SUBTLE BUGS IN MORE DEMANDING CIRCUMSTANCES."5.1 The end-to-end argument"JUST BECAUSE AN APPLICATION USES A DATA SYSTEM THAT PROVIDES COMPARATIVELY STRONG SAFETY PROPERTIES, SUCH AS SERIALIZABLE TRANSACTIONS, THAT DOES NOT MEAN THE APPLICATION IS GUARANTEED TO BE FREE FROM DATA LOSS OR CORRUPTION. If an application has a bug that causes it to write incorrect data or delete data, SERIALIZABLE TRANSACTIONS AREN'T GOING TO SAVE YOU.""THE FUNCTION IN QUESTION CAN COMPLETELY AND CORRECTLY BE IMPLEMENTED ONLY WITH THE KNOWLEDGE AND HELP OF THE APPLICATION STANDING AT THE ENDPOINTS OF THE COMMUNICATION SYSTEM. THEREFORE, PROVIDING THAT QUESTIONED FUNCTION AS A FEATURE OF THE COMMUNICATION SYSTEM ITSELF IS NOT POSSIBLE. (Sometimes an incomplete version provided by the communication system MAY BE USEFUL AS A PERFORMANCE ENHANCEMENT.)"5.2 Enforcing constraints without distributed transactions5.3 Timeliness vs Integrity — the chapter's most important distinction"AN INTERESTING PROPERTY OF EVENT-BASED DATAFLOW SYSTEMS IS THAT THEY DECOUPLE TIMELINESS AND INTEGRITY. When processing asynchronously, THERE IS NO GUARANTEE OF TIMELINESS unless you explicitly build consumers that wait. HOWEVER, INTEGRITY IS IN FACT CENTRAL TO STREAMING SYSTEMS. Exactly-once semantics IS A MECHANISM FOR PRESERVING INTEGRITY."5.4 Loosely interpreted constraints — the business-reality argument"IF THE COST OF THE APOLOGY IS ACCEPTABLE, THE TRADITIONAL MODEL OF CHECKING ALL CONSTRAINTS BEFORE EVEN WRITING THE DATA IS UNNECESSARILY RESTRICTIVE. It may well be reasonable to GO AHEAD WITH A WRITE OPTIMISTICALLY AND CHECK THE CONSTRAINT AFTER THE FACT. You can still ensure that VALIDATION OCCURS BEFORE TAKING ACTIONS THAT WOULD BE EXPENSIVE TO RECOVER FROM, BUT THAT DOESN'T IMPLY YOU MUST VALIDATE BEFORE YOU EVEN WRITE THE DATA."5.5 Coordination-avoiding data systemsTHE APOLOGY CALCULUS: "Coordination and constraints REDUCE THE NUMBER OF APOLOGIES YOU HAVE TO MAKE FOR INCONSISTENCIES, BUT POTENTIALLY ALSO REDUCE THE PERFORMANCE AND AVAILABILITY OF YOUR SYSTEM, AND THUS POTENTIALLY INCREASE THE NUMBER OF APOLOGIES YOU HAVE TO MAKE FOR OUTAGES. YOU CANNOT REDUCE THE NUMBER OF APOLOGIES TO ZERO, BUT YOU CAN AIM TO FIND THE BEST TRADE-OFF — THE SWEET SPOT WITH NEITHER TOO MANY INCONSISTENCIES NOR TOO MANY AVAILABILITY PROBLEMS."5.6 Trust, but Verify"ACID consistency is based on the idea that the database starts in a consistent state and a transaction transforms it to another consistent state. HOWEVER, THIS NOTION MAKES SENSE ONLY IF WE ASSUME THE TRANSACTION IS FREE FROM BUGS. IF THE APPLICATION USES THE DATABASE INCORRECTLY — FOR EXAMPLE, USING A WEAK ISOLATION LEVEL UNSAFELY — THE INTEGRITY OF THE DATABASE CANNOT BE GUARANTEED.""IF YOU WANT TO BE SURE THAT YOUR DATA IS STILL THERE, YOU HAVE TO READ IT AND CHECK. By the same argument, IT IS IMPORTANT TO TRY RESTORING FROM YOUR BACKUPS FROM TIME TO TIME — OTHERWISE YOU MAY FIND OUT THAT YOUR BACKUP IS BROKEN WHEN IT IS TOO LATE AND YOU HAVE ALREADY LOST DATA. DON'T JUST BLINDLY TRUST THAT IT IS ALL WORKING.""CHECKING THE INTEGRITY OF DATA SYSTEMS IS BEST DONE IN AN END-TO-END FASHION. THE MORE SYSTEMS WE CAN INCLUDE IN AN INTEGRITY CHECK, THE FEWER OPPORTUNITIES THERE ARE FOR CORRUPTION TO GO UNNOTICED. IF WE CAN CHECK THAT AN ENTIRE DERIVED DATA PIPELINE IS CORRECT END TO END, THEN ANY DISKS, NETWORKS, SERVICES, AND ALGORITHMS ALONG THE PATH ARE IMPLICITLY INCLUDED IN THE CHECK."