Skip to main content
Version: V2-Next

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.

ComponentResponsibility
OrchestratorRuns the workflow, keeps the state, moves data between steps, compensates after a failure, reports the result
AdapterRuns one command and returns the result

The four workflows​

ProcessSaga typePurpose
dataset-create.bpmnCREATEProvision a Dataset
dataset-update.bpmnUPDATEApply a change to a provisioned Dataset
dataset-unrelease.bpmnUNRELEASEWithdraw a released Dataset
dataset-delete.bpmnDELETERemove a Dataset

The create workflow​

The process runs these steps. A gateway before a step tests whether the Dataset needs it.

GatewayStepAdapter
Has FROST sink?Create FROST projectFROST
Create APISIX routeAPISIX
Has geo sink?Provision PostGIS sinkPostGIS
Create GeoServer workspaceGeoServer
Create GeoServer datastoreGeoServer
Has layers?Provision GeoServer layersGeoServer
Has pipelines?Deploy pipelinesNiFi
Publish success result—

Messages​

The orchestrator uses two Kafka topics.

TopicDirectionContent
de.civitascore.dataset.saga.triggerPortal backend → orchestratorA typed trigger: DatasetCreate, DatasetUpdate, DatasetUnrelease or DatasetDelete
de.civitascore.saga.resultOrchestrator → portal backendSAGA_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 compensationErrors and 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.