Skip to main content
Version: 2.0.0

Message Bus

The platform uses Apache Kafka as its central message bus. Components communicate asynchronously over the bus. ADR 026 records that decision.

warning

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​

CriterionAssessment
ScalingKafkaApache 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 managementA topic structure with fine granularity is possible.
AuthorizationKafkaApache 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.
EcosystemKafkaApache 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.
CommunityThe open-source community is large, which gives long-term stability.
PayloadKafkaApache 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.

→ Saga pattern

Topic names​

A topic name encodes the domain, the entity and the event. The structure makes access control and routing possible.

→ Topic configuration

Technology​

PartTechnology
Message busApache 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
EnvelopeCloudEvents
Adapter SDKPlain Java with few dependencies (ADR 016)
SerializationJackson, JSON
Saga stateThe 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.
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 013Use cloudevents Standard for Bus based Configuration Communication
ADR 016Configuration 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 021Definition of Configuration Events
ADR 026Select Message Bus
ADR 030Asynchronous Outbox for Config Adapter Synchronization
ADR 031Orchestrated 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 036Event driven Communication and loose coupling
ADR 041Revised Topic Naming Convention for Configuration Events
ADR 043KafkaApache 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