State Management
One of the most important considerations when designing an application or system composed of a number of independent components is how to handle state. There are often competing goals that govern the approach that is taken
Performance
we want to minimise the overhead of handling data as it passes through the system.
Safety
we want to ensure that data is not lost during normal operation of the system
Robustness
we want to ensure that if a component were to fail, we can minimise the loss of service and again ensure that any data loss is minimised.
Normally the steps taken to optimise safety and robustness have costs that can adversely affect performance. Finding the right balance across these goals is therefore critical to building a system that is both fast and reliable.
This section discusses the ways in which state can be managed in Chronicle Services applications.
Microservices and State
The Microservices pattern aims to decouple services from each other, in order to maximise the flexibility and resilience of the overall application. A key part of this is to isolate state to the microservice that requires it, and will modify it. Sharing of data between components can be achieved by passing relevant details between them in messages, or whatever communications mechanism that is being used.
Much of the state that is active in a microservice is transient, in that it is primarily relevant for processing a single message or request. However, normally there will be data whose value should be preserved for longer than a single request; for example cumulative values such as maxima or minima, or collections to which data is added when processing incoming messages.
Given that a further goal of the Microservices pattern is to enable a component to be stopped and restarted, perhaps after an outage, and resume processing with the same accumulated state, it is necessary to provide some means of preserving elements of state so that they can be restored. This has traditionally involved the use of a database of some kind, in most cases a Relational Database.
However, it has become clear that traditional relational databases have limitations when operating in a distributed environment. At the very least they introduce a significant overhead as a result of writing state changes to persistent storage.
One approach to minimising this overhead is to design and build components so that they have no state that requires to be persisted. This leads to components often described as "stateless".
Designing a stateless component requires incoming requests to carry the values that would otherwise have been implemented as persistent state, which can complicate the communications protocols used between components. There are also some components that will not practically be able to run in a stateless manner, so alternative approaches are needed.
Managing State Through Events
In order to deal with the shortcomings of traditional databases, a number of alternative models have evolved for managing state in distributed applications. An increasingly popular approach is to use events as a means of communicating changes in state from one component to others.
Applications that utilise this approach will typically follow the guidelines of Event-Driven Architecture, which require events to be persisted indefinitely. Any changes to key application state are notified by the relevant component using some form of event transport, which will persist the event and notify other components that register an interest.
The current value of an element of state is the result of applying all events that notify a change in its value. This is a radical change from the traditional way of managing state, where a change is represented by overwriting a single value and somehow persisting that change.
There are some important advantages to the event-driven approach, not least that we have the ability to examine the value of state at a given point in time, not just its current value. Essentially we have a built-in audit trail.
This is not to say that there are no disadvantages in using events like this, especially when we have a need to work with the "current value". Rather than replay all change events every time this is required, a component will normally cache the current value. When the component starts, or restarts, it is possible to build the cached value from state change events that have occurred before.
The event-based approach to managing state in applications has evolved into a set of patterns often referred to as "Event Sourcing", and several sophisticated event management systems have evolved to support it.
The Chronicle Services Approach to State
The basic model of a Chronicle Services application is of a number of independent processing components, interacting with each other using events that are passed using Chronicle Queue instances.
Chronicle Queue is a "store everything" data structure, in other words all events posted by a service will be retained in perpetuity on persistent storage. As such, it is well suited to following the event-based approach to managing state described above.
Where non-transient state is required in a service, Chronicle Services allows changes to this state to be represented in a Chronicle Queue instance’s persistent storage using events. Where necessary, the current value can be recreated by applying these events.
Service Startup and Event Replay
When a Chronicle Services component is started, or restarted, there will normally be a requirement to construct (or reconstruct) state so that it has the same values as when the component was stopped.
Since we have all the input events that were processed by the service to construct the state that was current when the service stopped, we can recreate this state by replaying all of these events through the appropriate event handlers. The events are read from the queues and replayed in order, but not in "real time" (i.e. without the delays between the events). Additionally, no output events will be posted from the handlers when the state is updated during a replay, as this would cause inconsistencies in the operation of downstream services.
Chronicle Service provides a number of ways in which event replay can be performed, known as Restart Strategies. These are managed through the service’s configuration, in particular using two configuration parameters:
startFromStrategy
Defines the point at which the services begins processing input events
inputsReplayStrategy
Defines from where events used to construct state are read
Default behaviour is that no state construction is performed, and input events are processed starting from the point at which the last output event was posted (in the case of a service starting for the first time, this will cause processing to begin with the first input event).
To read more about Event Replay in Chronicle Services please head to our portal.