Skip to main content
Version: 2.0.0

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.

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 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.bpmnUPDATEApply 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.bpmnUNRELEASEWithdraw 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.bpmnDELETERemove 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.

GatewayStepAdapter
Has FROST sink?Create FROST projectFROST
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. routeAPISIXApache 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 sinkPostGIS
Create GeoServer workspaceGeoServer
Create GeoServer datastoreGeoServer
Has layers?Provision GeoServer layersGeoServer
Has pipelines?Deploy pipelinesNiFiApache 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.

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 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 compensationErrors and 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.

  • 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