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 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
| Criterion | Assessment |
|---|---|
| Scaling | Kafka scales horizontally for a high message volume. |
| Topic management | A topic structure with fine granularity is possible. |
| Authorization | Kafka controls read and write access per topic. |
| Ecosystem | Kafka has an SDK for many languages and frameworks. |
| Community | The open-source community is large, which gives long-term stability. |
| Payload | Kafka 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.
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 Kafka, 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
| ADR | Title |
|---|---|
| ADR 013 | Use cloudevents Standard for Bus based Configuration Communication |
| ADR 016 | Configuration Adapter (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 Saga for Multi-Adapter Provisioning |
| ADR 036 | Event driven Communication and loose coupling |
| ADR 041 | Revised Topic Naming Convention for Configuration Events |
| ADR 043 | Kafka Authentication via SASL/SCRAM-SHA-512 |