Message Bus
The platform uses Apache Kafka as its central message bus. Components communicate asynchronously over the bus. ADR 026 records that decision.
Parts of this design are not yet built. The Transactional Outbox of ADR 030 has no implementation, and the installed 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. cluster has no authentication and no ACL.
Why a message bus
The platform is a set of services that deploy independently. Direct synchronous calls between them limit resilience and scaling. The architecture therefore applies three principles:
- Asynchronous first. A state change goes over the bus. A synchronous call is only for an operation that needs an immediate answer, for example a query from a user.
- Loose coupling. A component publishes an event. It does not know which component reads the event.
- Least privilege. A component gets access only to the topics it needs.
Why Kafka
| Criterion | Assessment |
|---|---|
| Scaling | 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. scales horizontally for a high message volume. |
| Topic management | A topic structure with fine granularity is possible. |
| Authorization | 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. controls read and write access per topic. |
| Ecosystem | 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. has an SDK for many languages and frameworks. |
| Community | The open-source community is large, which gives long-term stability. |
| Payload | 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 a small configuration event and a large data payload. |
Two categories of message
The bus carries two categories of message. They differ in purpose and in processing.
Configuration messages carry a management operation, for example the creation of a user or of a route. A 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. reads the message and applies the change to an external system. The envelope is CloudEvents (ADR 013).
Payload data carries domain content, for example a sensor reading. A Pipeline produces and consumes this data.
The CloudEvents envelope
Every configuration event uses the CloudEvents envelope. The envelope gives four advantages:
- The metadata structure is the same for every event type.
- The semantics do not depend on the transport.
- 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. maps the envelope to message headers, so a consumer can route on the event type without a deserialization of the payload.
- The payload schema can differ per event type.
The platform defines its own event types on the standard. ADR 021 gives the definitions.
Patterns
Transactional Outbox
ADR 030 decided this pattern: the entity change and the event go into one database transaction, and a separate process publishes the event to 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.. The platform does not have it. The portal backend still publishes inside the transaction.
Orchestrated Saga
Provisioning 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. needs steps in 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., and the steps must run in sequence. An orchestrator runs the workflow, keeps the state and rolls back the completed steps if a later step fails.
Topic names
A topic name encodes the domain, the entity and the event. The structure makes access control and routing possible.
Technology
| Part | Technology |
|---|---|
| Message bus | Apache 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., installed with Strimzi |
| Envelope | CloudEvents |
| Adapter SDK | Plain Java with few dependencies (ADR 016) |
| Serialization | Jackson, JSON |
| Saga state | The Flowable schema in PostgreSQL |
Design constraints
- Event versions need governance. The number of event types grows, so a schema rule is necessary.
- A consumer must be idempotent. A retry must not change the result.
- The system is eventually consistent. The user interface and the operating procedures must show this.
- Operation needs observability. Event tracing, a dead-letter path and monitoring of the synchronization state are necessary.
Related decisions
| ADRArchitecture Decision RecordA document capturing an architecture decision. Each ADR has a stable identifier and short title and is managed through four lifecycle states: Proposed, Accepted, Deprecated and Superseded. | Title |
|---|---|
| ADR 013 | Use cloudevents Standard for Bus based Configuration Communication |
| ADR 016 | 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. (CA) SDK Language |
| ADR 021 | Definition of Configuration Events |
| ADR 026 | Select Message Bus |
| ADR 030 | Asynchronous Outbox for Config Adapter Synchronization |
| 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 036 | Event driven Communication and loose coupling |
| ADR 041 | Revised Topic Naming Convention for Configuration Events |
| ADR 043 | 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. Authentication via SASLSASL (Simple Authentication and Security Layer)A framework defined in RFC 4422 that decouples authentication from application protocols./SCRAMSCRAM (Salted Challenge Response Authentication Mechanism)A challenge-response authentication protocol defined in RFC 5802.-SHA-512 |