Queue support
Besides its RESTful APIs, openFHIR Engine can be driven by a message broker: it subscribes to queues carrying FHIR payloads and/or openEHR compositions, translates every message and publishes the result to another queue. An integration that already moves its data over a broker does not need a REST-calling consumer in front of the engine.
Note
This is an enterprise-only feature and needs the messaging license option. Enabling it with a license that does
not carry the option stops the engine at startup with a message saying so. Licenses issued before this feature existed
do not carry it and have to be re-issued — contact us at info@open-fhir.com.
Apache Kafka is the supported broker.
How it works
For every message of a subscribed topic the engine
decides whether the message is FHIR or openEHR (see Payload detection),
translates it under the subscription’s tenant, with the same mappers, templates and terminology the REST operations use — a FHIR message is mapped as
$toopenehrmaps it, an openEHR message as$tofhirdoes,publishes the result to the subscription’s target topic,
and only then commits the message as consumed.
A message that cannot be translated is published, unchanged, to the subscription’s dead-letter topic together with an
OperationOutcome saying why. See Failures and delivery guarantees.
Configuration
Queue support is disabled by default. It is configured under openfhir.messaging:
openfhir:
messaging:
enabled: true
kafka:
bootstrap-servers: localhost:9092
properties: {} # shared client properties (security.protocol, sasl.*, ssl.*)
consumer: {} # consumer-only overrides
producer: {} # producer-only overrides (e.g. compression.type, max.request.size)
retry:
max-attempts: 5
initial-interval-ms: 1000
multiplier: 2
max-interval-ms: 30000
subscriptions:
- id: mixed
tenant: "123"
source: openfhir.in
payload: auto # auto | fhir | openehr
target: openfhir.out # or target-fhir / target-openehr
error-target: openfhir.dlq
template-id: # optional default, the header wins
openehr-format: canonical # canonical | flat
provenance: true
concurrency: 1
group-id: openfhir-mixed # default openfhir-<id>
Environment variable |
Default |
Description |
|---|---|---|
|
|
Enables queue support. When |
|
|
Kafka bootstrap servers, |
|
— |
Any Kafka client property shared by the consumers and the producer, e.g.
|
|
— |
Consumer-only Kafka properties, e.g. |
|
— |
Producer-only Kafka properties, e.g. |
|
|
How many times a message is delivered before it is dead-lettered, the first delivery included. Applies to transient failures only. |
|
|
Wait before the first retry. |
|
|
Factor the wait grows by with every retry. |
|
|
Upper bound of the wait. Must stay below the consumer’s |
Subscriptions are a list; as environment variables they are indexed, OPENFHIR_MESSAGING_SUBSCRIPTIONS_0_ID,
OPENFHIR_MESSAGING_SUBSCRIPTIONS_0_SOURCE and so on:
Subscription key |
Default |
Description |
|---|---|---|
|
— |
Required and unique. Names the subscription in the log and in |
|
— |
Required. The tenant every message of this subscription is translated under (see Multitenancy). A producer cannot choose or override it. |
|
— |
Required. The topic to consume. |
|
|
What the source topic carries: |
|
— |
Topic the results are published to, whatever their kind. |
|
— |
Topic for the produced FHIR Bundles / openEHR compositions; each wins over |
|
— |
Dead-letter topic. When unset, a message that cannot be translated is logged (without its payload) and dropped; a warning at startup points this out. |
|
— |
Template id used for messages that carry no |
|
|
Format of the produced compositions, |
|
|
Whether produced Bundles carry the engine-generated |
|
|
Consumer threads of this subscription on one node. Effective parallelism is |
|
|
Kafka consumer group. All nodes of a cluster share it. |
One mixed queue or one queue per kind
A topic that carries both FHIR and openEHR is one subscription with payload: auto, as in the example above. One
topic per kind is two subscriptions:
subscriptions:
- id: compositions
tenant: "123"
source: ehr.compositions
payload: openehr
target-fhir: fhir.bundles
error-target: openfhir.dlq
- id: bundles
tenant: "123"
source: fhir.incoming
payload: fhir
target-openehr: ehr.incoming
error-target: openfhir.dlq
The configuration is validated at startup and the engine does not start when it is invalid: subscription ids must be
unique, every direction a subscription can produce needs a target, and no subscription may publish to a topic that any
subscription consumes — the engine would otherwise consume its own output, and with payload: auto translate it back
and forth forever.
Securing the connection
Everything the Kafka clients understand can be set under kafka.properties, for example SASL over TLS:
openfhir:
messaging:
kafka:
bootstrap-servers: broker-1:9093,broker-2:9093
properties:
security.protocol: SASL_SSL
sasl.mechanism: SCRAM-SHA-512
sasl.jaas.config: org.apache.kafka.common.security.scram.ScramLoginModule required username="openfhir" password="...";
ssl.truststore.location: /app/truststore.jks
ssl.truststore.password: "..."
GET /status masks these values (see RESTful APIs). spring.kafka.* properties are not read; openfhir.messaging.kafka
is the only place Kafka is configured.
A few client settings guard the delivery guarantee and cannot be overridden: the consumers never auto-commit, and the
producer always runs with acks=all and idempotence enabled. Two consumer defaults differ from Kafka’s own and can be
overridden under kafka.consumer: auto.offset.reset is earliest (a new subscription starts at the beginning of
its topic instead of skipping what is already there) and max.poll.records is 10.
Message contract
The body of a message is the payload itself — a FHIR Bundle (or single resource) or an openEHR composition (canonical or flat) as JSON. There is no envelope: everything else travels in headers.
Inbound headers
All of them are optional.
Header |
Effect |
|---|---|
|
|
|
Template id to map with. Required for a flat composition, as it is for |
|
|
|
Correlation id, echoed on the result and recorded on the mapping insight. Defaults to
|
|
The call context of |
Outbound headers
Header |
Meaning |
|---|---|
|
|
|
On a target topic |
|
An |
|
The total number of issues, also when |
|
The request id, echoed or generated. |
|
The template the message was mapped with. |
|
Where the translated message came from. Together they identify it uniquely — use them to detect duplicates. |
The record key of a message is preserved on its result and on its dead-letter record, so related messages keep landing on the same partition.
The body of a result is the bare FHIR Bundle or the bare composition. Unlike the REST operations, issues are not put
into the body (no OperationOutcome Bundle entry, no Parameters wrapper), so a consumer can use it as it is.
How much openfhir-outcome carries follows OPENFHIR_OPERATIONS_OUTCOME_VERBOSITY (see RESTful APIs); with
none neither it nor openfhir-issue-count is set, while openfhir-status still tells partial from ok.
A dead-letter record carries the original body, key and headers unchanged — it can be put back on the source topic as it is — with the outbound headers added.
Payload detection
A body with a top-level resourceType is FHIR. A body with a top-level _type (canonical) or with flat paths as
member names is openEHR. For an array, its first element decides.
The body is looked at in every mode. A subscription’s payload setting and the openfhir-payload-type header say
what is expected; a body that contradicts either of them — a FHIR Bundle on a payload: openehr subscription, say —
is dead-lettered rather than mapped as something it is not. A body that is JSON but says nothing about itself needs the
header or a non-auto subscription. A body that is not JSON is dead-lettered.
Failures and delivery guarantees
What happened |
What the engine does |
|---|---|
A single mapping failed, or elements were skipped |
The result is published with |
The message is at fault: not JSON, unknown template, no context mapper, missing template id, … |
Dead-lettered straight away, without retrying. These are the failures REST answers with a |
Something the engine depends on is unavailable: database, CDR, the broker itself |
Retried with exponential backoff ( |
The dead-letter topic cannot be written |
The message is retried until it can: the partition waits, nothing is lost. Only a record the dead-letter topic rejects for its size is logged and skipped. |
In every case the subscription keeps consuming afterwards.
Delivery is at-least-once. A message is committed only after its result (or its dead-letter record) was acknowledged
by the broker, so a crash or a rebalance in between means the message is translated and published again. Produced
Bundles contain generated ids, so two results of the same message are not byte-identical: consumers that need
exactly-once semantics should dedupe on the openfhir-source-topic / -partition / -offset headers.
Order is preserved per partition, and therefore per record key.
Note
A failure of infrastructure inside a single mapping — a terminology server that does not answer, for example —
surfaces as an error issue on a partial result, as it does over REST. It is not retried.
Startup, shutdown and clustering
Subscriptions start once the engine is fully up: after the license was validated and the
bootstrap scan ran. A broker that cannot be reached never keeps the engine from starting or
from serving REST: the subscription is reported as FAILED and started again every 30 seconds until it succeeds. A
source topic that does not exist yet is waited for.
On shutdown the subscriptions are stopped first and the producer is closed after them, so a message in flight can still publish its result.
In a multi-node deployment every node runs the same subscriptions in the same consumer group; Kafka splits the partitions of a topic between the nodes and moves them when a node leaves. No further coordination is needed. A rebalance can deliver a message a second time (see above).
Monitoring
GET /status reports each subscription’s state and what happened to the messages this node took since it started:
{
"messaging": {
"enabled": true,
"subscriptions": [
{
"id": "mixed",
"tenant": "123",
"source": "openfhir.in",
"payload": "auto",
"state": "RUNNING",
"processed": 1284,
"partial": 3,
"deadLettered": 1
}
]
}
}
processed counts messages published to a target topic, partial those among them that carried a mapping error,
deadLettered those that could not be translated (including the ones dropped for want of an error-target).
state is one of STARTING, RUNNING, FAILED and STOPPED; RUNNING means the consumer is running, not
that the broker currently answers — a consumer that lost its broker keeps reconnecting in the background.
Every message is also recorded as a mapping insight, with its request id, direction and — for a dead-lettered message — the error.
Sizing and limits
Results are bounded by the message size limits of the broker and the producer, 1 MB by default. FHIR Bundles are larger than the compositions they are produced from: consider
kafka.producer.compression.type: gzipand raisemax.request.size(and the topic’smax.message.bytes) where needed. A result that is still too large is dead-lettered with anOperationOutcomesaying so.One mapping insight is written per message. On a high-volume topic consider
OPENFHIR_INSIGHTS_ENABLED=false, or at least leave payload storage off.Retries wait on the consumer thread, so a message that is being retried holds back the messages behind it on the same partition; other partitions are unaffected.