Orchestrated Saga Pattern
Problem
To provision a DatasetDatasetA data-related element that contains processed data and makes it available for consumption. A Dataset is populated via Pipelines and carries Metadata and access permissions., the platform must call more than one Configuration AdapterConfiguration AdapterThe component that consumes data models on behalf of platform components that cannot consume them directly, and configures the component accordingly. In the secrets management flow, the Configuration Adapter is the sole component that resolves Vault references into concrete credentials.. 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 AdapterConfiguration AdapterThe component that consumes data models on behalf of platform components that cannot consume them directly, and configures the component accordingly. In the secrets management flow, the Configuration Adapter is the sole component that resolves Vault references into concrete credentials. 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 AdapterConfiguration AdapterThe component that consumes data models on behalf of platform components that cannot consume them directly, and configures the component accordingly. In the secrets management flow, the Configuration Adapter is the sole component that resolves Vault references into concrete credentials. 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 DatasetDatasetA data-related element that contains processed data and makes it available for consumption. A Dataset is populated via Pipelines and carries Metadata and access permissions. |
dataset-update.bpmn | UPDATE | Apply a change to a provisioned DatasetDatasetA data-related element that contains processed data and makes it available for consumption. A Dataset is populated via Pipelines and carries Metadata and access permissions. |
dataset-unrelease.bpmn | UNRELEASE | Withdraw a released DatasetDatasetA data-related element that contains processed data and makes it available for consumption. A Dataset is populated via Pipelines and carries Metadata and access permissions. |
dataset-delete.bpmn | DELETE | Remove a DatasetDatasetA data-related element that contains processed data and makes it available for consumption. A Dataset is populated via Pipelines and carries Metadata and access permissions. |
The create workflow
The process runs these steps. A gateway before a step tests whether the DatasetDatasetA data-related element that contains processed data and makes it available for consumption. A Dataset is populated via Pipelines and carries Metadata and access permissions. needs it.
| Gateway | Step | Adapter |
|---|---|---|
| Has FROST sink? | Create FROST project | FROST |
| Create APISIXApache APISIXAn open-source API gateway for traffic management, security and observability. In CIVITAS/CORE it is used as the centralized entrypoint to route and protect externally exposed APIs. route | APISIXApache APISIXAn open-source API gateway for traffic management, security and observability. In CIVITAS/CORE it is used as the centralized entrypoint to route and protect externally exposed APIs. | |
| 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 | NiFiApache NiFiA stream processing and connector framework. In CIVITAS/CORE it is the pipeline engine for data integration and transformation (see ADR 047) and implements dataset-defined data flows. |
| Publish success result | — |
Messages
The orchestrator uses two KafkaApache KafkaA distributed event streaming platform. In CIVITAS/CORE it is used as the message bus to transport events, models and data in data flows. 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 KafkaApache KafkaA distributed event streaming platform. In CIVITAS/CORE it is used as the message bus to transport events, models and data in data flows. message. The orchestrator calls the handler of the adapter directly, through a registry that holds one handler per adapter name. KafkaApache KafkaA distributed event streaming platform. In CIVITAS/CORE it is used as the message bus to transport events, models and data in data flows. carries only the start of the saga and the final result.
The portal backend records a saga in flight on the DatasetDatasetA data-related element that contains processed data and makes it available for consumption. A Dataset is populated via Pipelines and carries Metadata and access permissions.. The column
pending_saga_type holds CREATE, UPDATE, UNRELEASE or DELETE. A second operation
on the same DatasetDatasetA data-related element that contains processed data and makes it available for consumption. A Dataset is populated via Pipelines and carries Metadata and access permissions. 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 APISIXApache APISIXAn open-source API gateway for traffic management, security and observability. In CIVITAS/CORE it is used as the centralized entrypoint to route and protect externally exposed APIs. route, the PostGIS sink, the GeoServer workspace and the NiFiApache NiFiA stream processing and connector framework. In CIVITAS/CORE it is the pipeline engine for data integration and transformation (see ADR 047) and implements dataset-defined data flows. 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 KafkaApache KafkaA distributed event streaming platform. In CIVITAS/CORE it is used as the message bus to transport events, models and data in data flows. 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 SagaSaga PatternA pattern for coordinating long-running, multi-step operations across several components without a distributed transaction. Each step has a compensating action; if a later step fails, the compensating actions of all previously completed steps are executed in reverse order. In CIVITAS/CORE it is used for provisioning workflows that require sequential, cross-adapter operations. for Multi-Adapter Provisioning
- ADR 013: Use cloudevents Standard for Bus based Configuration Communication