Skip to main content
Version: V2-Next

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 Kafka 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
ScalingKafka scales horizontally for a high message volume.
Topic managementA topic structure with fine granularity is possible.
AuthorizationKafka controls read and write access per topic.
EcosystemKafka has an SDK for many languages and frameworks.
CommunityThe open-source community is large, which gives long-term stability.
PayloadKafka 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 Adapter 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.
  • Kafka 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 Kafka. The platform does not have it. The portal backend still publishes inside the transaction.

Orchestrated Saga​

Provisioning a Dataset needs steps in more than one Configuration Adapter, 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 Kafka, 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.
ADRTitle
ADR 013Use cloudevents Standard for Bus based Configuration Communication
ADR 016Configuration Adapter (CA) SDK Language
ADR 021Definition of Configuration Events
ADR 026Select Message Bus
ADR 030Asynchronous Outbox for Config Adapter Synchronization
ADR 031Orchestrated Saga for Multi-Adapter Provisioning
ADR 036Event driven Communication and loose coupling
ADR 041Revised Topic Naming Convention for Configuration Events
ADR 043Kafka Authentication via SASL/SCRAM-SHA-512