How to get transactions between almost any data stores
petereliaskraft.net
petereliaskraft.net
I know it’s canon at this point, but please stop using this example people. No bank ever worked this way, and often outsiders who hear about this example get the wrong idea.
This is really more like a CRDT than like transaction control.
It's also where Bitcoin took its inspiration from, the blockchain is just one big ledger with some hashing on top.
Or, are you old enough to have learned to reconcile a checkbook? Where you keep track of all the checks you write and what your balance should be, and then mark them as confirmed when they show up on your (monthly, because paper via usps) bank statements?
Or the queue could be strongly ordered.
I am curious how this kind of thing can scale. What I'm inspired by this article is the idea of applying multiversion semantics to datastores that aren't multiversioned.
Sounds like if you have a transactionally sound store of "versions" and a way to query every datastore with a version field, you could implement point in time consistency, because the current version would only change when all datastores are updated. You still have to wait for all datastores to be updated, so that part is similar to 2 phase commit.
My understanding is that two phase commit does not scale and is slow because of all the round trips, but we still want linearizability and strong consistency in systems.
I tried to implement an asynchronously replicating data protocol but it is eventually consistent, not linearizable (it fails the Jepsen linearizability test)
https://github.com/samsquire/eventually-consistent-mesh
I'm a beginner in this area but I have a toy multiversion concurrency control implementation that uses Java threads. (See MVCC.java and TransactionC.java) the key statements of code is in commit() and shouldRestart().
https://github.com/samsquire/multiversion-concurrency-contro...
Dropbox does something 2 phase commit with their cross shard replication which is interesting https://dropbox.tech/infrastructure/cross-shard-transactions...
SQLite can achieve this with some help: https://blog.expensify.com/2018/01/08/scaling-sqlite-to-4m-q...
More seriously, I suspect that plain SQLite by itself, running in process in the main service, writing to a ':memory:' database for ultimate speed or perhaps to a local file with WAL mode for slightly more durability, could achieve pretty nice efficiency and work as the main transaction manager database.
I did something like this a while ago; some operations in my system inserted multiple records into different collections (into a database which did not support atomic transactions) and I needed to guarantee that when a record existed in a particular collection in the 'settled' state that it would guarantee the existence of a matching record in a different collection. It was very reliable. I believe I didn't even end up needing to launch a separate process for the settlement logic since I was using Node.js, I just used an async task scheduled to run on an interval. The front end would just ignore records until they were in the settled state (but you could also show them as pending).
Another thing which really helps is using UUIDs are IDs as it guarantees idempotence so you don't have to worry about double-insertion in case of process failure/restart. You can just re-process some records that you may have already processed with no side effects; you just need to check that the error is an insertion conflict and continue.
Because I used UUID as IDs and they were created on the front end, the user in my app could potentially re-submit the form (click submit again after seeing the 'Unexpected connection error please try again...' message); this represented a second chance to complete the transaction. Those records which were already successfully inserted into the db the last time would be ignored the second time (due to ID conflict) and those which had failed the first time and not been inserted into the db would then be inserted as pending; then the settlement script could complete its job and settle on the next interval.
For example, if you settle based on creation timestamps then you know that if a record is settled, then all records with a smaller timestamp within the same collection/table are settled too. You can also implement a similar guarantee if dealing with multiple collections/tables.
For example, if you create a 'Company' record along with a 'User' record which points to the 'companyId' and you want them to be settled atomically as a pair, then in your settlement process/logic, you could ensure that you always mark the company as settled first (and wait to get a success response from that DB query) before you then mark the user as settled - In this case, you only need to check that the user is settled to know that the associated company is also settled since your settlement logic guarantees that it's not possible for the user record to be settled before its associated company since you check that the company settlement query was successful before moving onto the related user.
In this case, you can think of settlement as having two legs; the first leg is the company record, the second leg is the user record. If the second leg of the settlement fails, the user will stay in a pending state which may cause your settlement script to reprocess (or re-check) the first leg (the related company record) since, if the last leg of the settlement failed, it will treat it as if the entire transaction failed and that's fine.
In practice transactions between arbitrary data stores would result in potentially boundless and unpredictable latency, no?
Also, is Postgres strongly consistent and linearizable)? One alternative would be using a database with stronger consistency guarantees (Spanner is but not open source, FoundationDB is but has limitations on transactions unless you implement mvcc yourself, which to be fair you are).
There is still a concept of (transaction) isolation levels, and the ANSI SQL standard defines a transaction mode READ UNCOMMITTED that could give you inconsistent results, but Postgres ignores that and treats it as READ COMMITTED.
Replication speed could be bad, I don't see a reason to expect that.
PostgreSQL can have serializable transactions though that is not the default isolation level.
I wonder what happens for records with many changes during their lifetime. Wouldn't the size of such a snapshot grow infinetely for them?
IBM had CICS transactions back in the day for some definition of heterogeneity. Tuxedo did this in the 80s across broader platforms. In the 90s (when I heard of it) there was an open standard created for it: https://en.wikipedia.org/wiki/X/Open_XA so that Unix-y type systems had some agreement on how to do it without buying mainframe/tuxedo-y type components. And around or a bit after then, Microsoft created its own "Distributed Transaction Coordinator". This is not my area of expertise but I've heard about it throughout the last 30 years.
The NoSQL side of heterogeneous is newer than the 90s but even there, I found in 2 minutes a paper from 2006, 17 years ago, on the topic: https://www.researchgate.net/profile/Akon-Dey/publication/28...
This is definitely a wheel that gets reinvented.
The Open Group has defined an industry-standard model for transactional work that allows changes made against unrelated resources to be part of a single global transaction. An example of this is changes to databases that are provided by two separate vendors. This model is called the X/Open Distributed Transaction Processing model.
Source: https://www.ibm.com/docs/en/i/7.1?topic=concepts-xa-transact...
Is Microsoft's MSDTC more single-vendor? Even that seems to support multi-vendor databases through its support of Open/XA:
When the DTC acts as an XA-compliant transaction manager, Oracle, IBM DB/2, Sybase, Informix, and other XA-compliant resource managers can participate in transactions that the DTC controls.
Source: https://learn.microsoft.com/en-us/previous-versions/windows/...
Diving into the history of which distributed transaction manager was the first to support multi-vendor database transactions is left as an exercise for the reader...
P.S. Some good Oracle docs on distributed transactions that gives you a flavor of how it works: https://docs.oracle.com/cd/A97630_01/java.920/a96654/xadistr...
I suspect industry knowledge/awareness of this stuff waned as web-generated transactions boomed and cloud platform vendors like AWS sold "eventual consistency" to developers and through them back to management and end users (like myself!) who think "Oh, that's weird, why didn't X update? Let me just refresh my browser... oh, there it is".
> The traditional solution is to use two-phase commit through a protocol like X/Open XA. However, while XA is supported by most big relational databases like Postgres and MySQL, it's not supported by popular newer data stores like MongoDB, Cassandra, or Elasticsearch
Cassandra has a distributed transaction based on paxos, and I think a raft version is in the works.
I admittedly skimmed the article, but the notion of a central coordinator for the transaction isn't really a scalable solution. In the stated use case, heterogeneous stores, it's kind of what you're stuck with to some degree.
But... why not use a purely distributed central store like Cassandra or Zookeeper (I think zoo is masterless?) rather than Postgres?
IMO this doesn't have a chance of standing up to a network partition like Aphyr would throw at it, but then again I didn't graduate from Stanford.
Another article you didn’t need to read because it talks about a problem you likely don’t have because the author can’t be bothered to come up with actual, real world examples.
as someone who hasn't implemented transactions in microservices before, I have only seen the bank transfer example and it seems adequate and easy to understand - I wouldn't know what limitations it has
It’s fake. As a result, you can’t actually evaluate the engineering trade offs of the proposed solution.
Here’s the thing: do you even need transactions (in the sense of ACID transactions that are resolved heuristically, as described in this article) between your microservices? Likely not.
And what are the implications of having to wait around for PGSQL and Mongo negotiate a transaction between a third data store? Probably blocking and latency and all the associated problems which will lead to … not use transaction across service boundaries.
Both parties misunderstanding the other is a common trope in movies and novels for a reason.
The unreliable messaging links (internal mail) mean that resilience against missed messages and guards against not making progress must be built into the business process instead of an infrastructure layer.
Common actions like Create, Delete, and some Updates touch all of these services. Automatic transactions between these would be fantastic. Especially on the billing side, some billing flows involve multiple trips between Stripe and our DB, batching the whole process inside a guaranteed reversable transaction would be lovely.
Of course we're not going to get that with Stripe and Cognito as they expose bespoke API's, but the idea still holds.
Further, those complex interactions between those various services will be much hard to debug and tweak when trying to provide support to end users. Can’t pay your bill because Cognito is down? Can’t grant access because Stripe is down? The data warehouse is doing an index build so you can’t do any transactions? Guess we should turn off the website between 2AM and 4AM ET to allow for data warehouse rebuilds.
There’s a reason we don’t do things this way.
I have had this exact use case in a problem I was solving where the user was waiting at the other end. Please do not presume that only the problems you have solved are worth solving or exist in the real world.
>> Can’t pay your bill because Cognito is down?
I mean, if your AuthN and AuthZ is down, what else would you do ?!!
>> There’s a reason we don’t do things this way.
The smugness in this line is astounding. I don't think I can express how wrong you are.
Minor miscommunication… from a write standpoint. Obviously if you can’t auth, you are in trouble. But if the billing needs to happen so you can restore a permission, then you need to hit a card (for example) and then somehow update your authz service.
In most cases, it’s not strict availability that’s likely your problem (eg. Full down, no responses) it’s going to be business logic or config issues impacting a small subset of TX.
And, honest, if I’m so wrong about this, where is the XA distributed transaction coordinator for Cognito?
I’d be shocked if Amazon ever implemented a single-point-of-failure distributed transaction coordinator and exposed that from their services.
Its true that most people don't need it at all, and I cringe at the thought of a second wave of micro-service muppets stitching it all together with something like this, but there are actually real-world examples. I'm not going to offer you any though, due to your belligerence in the other comments.
A bank transfer is a good example because it is easy to understand, regardless of your background.