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

  1. decides whether the message is FHIR or openEHR (see Payload detection),

  2. translates it under the subscription’s tenant, with the same mappers, templates and terminology the REST operations use — a FHIR message is mapped as $toopenehr maps it, an openEHR message as $tofhir does,

  3. publishes the result to the subscription’s target topic,

  4. 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

OPENFHIR_MESSAGING_ENABLED

false

Enables queue support. When false, no broker connection is ever attempted.

OPENFHIR_MESSAGING_KAFKA_BOOTSTRAP_SERVERS

localhost:9092

Kafka bootstrap servers, host:port[,host:port].

OPENFHIR_MESSAGING_KAFKA_PROPERTIES_*

—

Any Kafka client property shared by the consumers and the producer, e.g. OPENFHIR_MESSAGING_KAFKA_PROPERTIES_SECURITY_PROTOCOL for security.protocol.

OPENFHIR_MESSAGING_KAFKA_CONSUMER_*

—

Consumer-only Kafka properties, e.g. OPENFHIR_MESSAGING_KAFKA_CONSUMER_MAX_POLL_RECORDS.

OPENFHIR_MESSAGING_KAFKA_PRODUCER_*

—

Producer-only Kafka properties, e.g. OPENFHIR_MESSAGING_KAFKA_PRODUCER_COMPRESSION_TYPE.

OPENFHIR_MESSAGING_RETRY_MAX_ATTEMPTS

5

How many times a message is delivered before it is dead-lettered, the first delivery included. Applies to transient failures only.

OPENFHIR_MESSAGING_RETRY_INITIAL_INTERVAL_MS

1000

Wait before the first retry.

OPENFHIR_MESSAGING_RETRY_MULTIPLIER

2

Factor the wait grows by with every retry.

OPENFHIR_MESSAGING_RETRY_MAX_INTERVAL_MS

30000

Upper bound of the wait. Must stay below the consumer’s max.poll.interval.ms (5 minutes unless overridden).

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

id

—

Required and unique. Names the subscription in the log and in GET /status.

tenant

—

Required. The tenant every message of this subscription is translated under (see Multitenancy). A producer cannot choose or override it.

source

—

Required. The topic to consume.

payload

auto

What the source topic carries: auto (FHIR and openEHR mixed, detected per message), fhir or openehr.

target

—

Topic the results are published to, whatever their kind.

target-fhir / target-openehr

—

Topic for the produced FHIR Bundles / openEHR compositions; each wins over target. They are named after what is published: an openehr subscription produces FHIR and therefore needs target-fhir (or target).

error-target

—

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

—

Template id used for messages that carry no openfhir-template-id header. Leave it unset on a topic that carries more than one template.

openehr-format

canonical

Format of the produced compositions, canonical or flat, unless the message’s openfhir-format header says otherwise.

provenance

true

Whether produced Bundles carry the engine-generated Provenance entry, as $tofhir responses do. The openfhir-who / openfhir-on-behalf-of headers only have an effect through it.

concurrency

1

Consumer threads of this subscription on one node. Effective parallelism is min(concurrency, partitions).

group-id

openfhir-<id>

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

openfhir-payload-type

fhir or openehr. Declares what the body is, for a body that does not say so itself (see Payload detection).

openfhir-template-id

Template id to map with. Required for a flat composition, as it is for $tofhir, unless the subscription has a template-id.

openfhir-format

canonical or flat: the format of the composition produced from a FHIR message.

openfhir-request-id

Correlation id, echoed on the result and recorded on the mapping insight. Defaults to <subscription>:<topic>-<partition>@<offset>.

openfhir-ehr-id, openfhir-patient, openfhir-who, openfhir-on-behalf-of

The call context of $tofhir (ehr_id, patient, who, onBehalfOf): used to fill empty subject references and the Provenance of the produced Bundle.

Outbound headers

Header

Meaning

openfhir-direction

tofhir or toopenehr. Absent on a dead-letter record whose kind could not be determined.

openfhir-status

On a target topic ok, or partial when at least one mapping failed and the body is the result of the remaining ones. On the dead-letter topic error.

openfhir-outcome

An OperationOutcome (JSON) with the issues reported during mapping, or the reason a message was dead-lettered. Capped at 16 KB: errors come first, and a last issue says how many were left out.

openfhir-issue-count

The total number of issues, also when openfhir-outcome was capped.

openfhir-request-id

The request id, echoed or generated.

openfhir-template-id

The template the message was mapped with.

openfhir-source-topic, openfhir-source-partition, openfhir-source-offset

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 openfhir-status: partial (or ok for warnings only) and the issues in openfhir-outcome — exactly what the REST operations answer with 200 and an OperationOutcome.

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 4xx.

Something the engine depends on is unavailable: database, CDR, the broker itself

Retried with exponential backoff (retry.*). If it still fails after max-attempts deliveries, the message is dead-lettered.

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: gzip and raise max.request.size (and the topic’s max.message.bytes) where needed. A result that is still too large is dead-lettered with an OperationOutcome saying 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.