Orchestrated Saga Pattern
Problem
To provision a Dataset, the platform must call more than one Configuration Adapter. The calls must run in sequence, because a later step needs the output of an earlier step. If a later step fails, the platform must remove what the earlier steps created.
The Configuration Adapter framework handles an independent operation well. It cannot manage a dependency between two adapters.
Why orchestration
ADR 031 examined choreography and rejected it. With three or more sequential steps, choreography gives four problems:
- Each adapter must know its predecessor and its successor.
- A rollback runs backwards through the chain. If one compensation fails, the next compensation never starts.
- Each adapter collects routing logic that does not belong to its domain.
- The first event must carry the configuration for every adapter.
Solution
A central orchestrator runs the workflow. An adapter stays a command handler: it receives
a command, it runs the command, and it returns a result. An adapter does not know about a
saga, and it does not know whether a delete is a normal operation or a rollback.
The orchestrator is the config-adapter-flowable module. It runs inside the
Configuration Adapter application. The workflow engine is
Flowable, and each workflow is a BPMN process.
| Component | Responsibility |
|---|---|
| Orchestrator | Runs the workflow, keeps the state, moves data between steps, compensates after a failure, reports the result |
| Adapter | Runs one command and returns the result |
The four workflows
| Process | Saga type | Purpose |
|---|---|---|
dataset-create.bpmn | CREATE | Provision a Dataset |
dataset-update.bpmn | UPDATE | Apply a change to a provisioned Dataset |
dataset-unrelease.bpmn | UNRELEASE | Withdraw a released Dataset |
dataset-delete.bpmn | DELETE | Remove a Dataset |
The create workflow
The process runs these steps. A gateway before a step tests whether the Dataset needs it.
| Gateway | Step | Adapter |
|---|---|---|
| Has FROST sink? | Create FROST project | FROST |
| Create APISIX route | APISIX | |
| Has geo sink? | Provision PostGIS sink | PostGIS |
| Create GeoServer workspace | GeoServer | |
| Create GeoServer datastore | GeoServer | |
| Has layers? | Provision GeoServer layers | GeoServer |
| Has pipelines? | Deploy pipelines | NiFi |
| Publish success result | — |
Messages
The orchestrator uses two Kafka topics.
| Topic | Direction | Content |
|---|---|---|
de.civitascore.dataset.saga.trigger | Portal backend → orchestrator | A typed trigger: DatasetCreate, DatasetUpdate, DatasetUnrelease or DatasetDelete |
de.civitascore.saga.result | Orchestrator → portal backend | SAGA_COMPLETED or SAGA_FAILED |
A step is not a Kafka message. The orchestrator calls the handler of the adapter directly, through a registry that holds one handler per adapter name. Kafka carries only the start of the saga and the final result.
The portal backend records a saga in flight on the Dataset. The column
pending_saga_type holds CREATE, UPDATE, UNRELEASE or DELETE. A second operation
on the same Dataset gets a SagaInFlightException while a saga runs.
Compensation
If a step fails, the process leaves the normal path over an error boundary event. The error path then calls the compensation of each completed step.
Three properties are important:
- The compensation is a plain service task. The process models do not use the compensation events of BPMN 2.0.
- The compensation runs in sequence, not in parallel.
- The compensation is best effort. A failed compensation does not stop the path. The
process collects the errors in the variable
compensationErrorsand reports them at the end.
The create process has a compensation for the FROST project, the APISIX route, the PostGIS sink, the GeoServer workspace and the NiFi pipelines.
State
Flowable keeps the state of a running process in its own schema in PostgreSQL. The platform does not have a saga state table of its own.
A Kafka record gives the business key of the process instance. The key is the topic, the partition and the offset of the trigger record, so the same record always maps to the same process instance.
Related documentation
- Message bus
- Topic configuration
- ADR 031: Orchestrated Saga for Multi-Adapter Provisioning
- ADR 013: Use cloudevents Standard for Bus based Configuration Communication