ACID like communication between microservices pattern

Viewed 196

I've been working on microservices architecture application for the past few months and I am still trying to get used to the distributed nature. There is one pattern I've noticed multiple times that I am not sure what's the preferred way to handle it.

Let's say we have service A and service B and service C. Service A exposes an API where one of the methods depends on calling an API exposed by B to create a resource RB and also depends on an API exposed by C to create resource RC. So in a perfect world A, B, and C are all working fine but the use case I've noticed a few times is that either B or C can be down during the execution of the API logic exposed by A. Even more how it should be tackled when RB is created, C is down so RC cannot be created and we try to rollback the creation of RB by calling let's say /delete/ on a service B but during that time B went down as well. Now RB is created but in the end, it shouldn't since RC failed and the execution of the API logic of A should.

The same A, B & C may be 3 nodes inside a cluster environment trying to propagate data across the cluster when data is posted to one of the nodes.

Sorry for the long text, thanks.

3 Answers

This issue is decades old and there have been many different ways to solve it. The problem is that the type of distributed transaction management needed to actually implement what you are describing is difficult to get right and tends to lead to overly complicated solutions. This complexity is really the reason why things like EJB have gone the way of the dodo.

Over time things have evolved to the point that in most cases it is easier to have things be eventually consistent and offload retrying to different types of message queues etc. (as described by previous posters). Of course there are situations where you simply cannot be eventually consistent, but they are not hard to identify and are the minority.

Persist those error messages to a queue and have a background job retry the rollback.

The problem you have is the rest calls between services, doing it this way you are doing temporal coupling between the services, so if in that moment any of the services is down all the operation is going to fail and the worst thing is that you can have an inconsistency in your data as you said.

The best approach to handle failure in a distributed system is make a reactive system. Here is the link to the reactive manifesto.

https://www.reactivemanifesto.org/

What they say in sort is that if you want a resilient system you must use asynchronous message passing for service communication, the disadvantage of this is that you have to live with eventual consistency as @qujck said, but it brings more advantages than disadvantages.

In your use case, when you make a call to service A. It should create some record to track the operation with a waiting for B and C state. Then, it should send commands (messages) to service B and C using some kind of message broker like Kafka with event persistence to ensure none message will be lost. If any of this service is down there is no problem because the message will be stay in kafka until the service is up and eventually consume it.

When this happens each of the service will send an event (message) saying 'Im finish'. Service A will listen to this message and then it will update the state of the operation to 'waiting for B', or C and finally to 'completed' when both messages arrive.

If for any other reason service B or C couldn't fullfill the request, they will send an error message instead of a finish one and then service A will send a command to the other service requesting a rollback. And if this service is down, it doesn't matter because eventually it will go up, read the rollback command an executes it to ensure there isn't any inconsistency

Related