Cluster Manager Processors

 

Cluster Manager module is available to all other modules, so no specific dependency needs to be configured to use it.

Broker Consumer/Processor

 

Smile Camel broker is both a consumer and a processor, which means that it can participate in a route as a consumer <from..> or producer <to..> endpoint.

  • Example from URI: <from uri="smile:clustermgr/broker?topic=my-kafka-from-topic">

  • Description: Takes messages from indicated topic, wraps them in a Camel Exchange and sends them to the following route node.

  • Example to URI: <to uri="smile:clustermgr/broker?topic=my-kafka-to-topic&messageType=ca.uhn.fhir.broker.api.RawStringMessage&payloadType=java.lang.String">

  • Description: Sends received Camel Exchange message to defined Smile internal topic.

Broker Processor Parameters

The following parameters are available for the broker processor:

  • topic – (required) The name of the topic to publish to or subscribe from.
  • messageType – (optional) Specifies the message type class to use when sending messages. Must be a class that implements IMessage. Default is ca.uhn.fhir.broker.api.RawStringMessage.
  • payloadType – (optional) Specifies the payload type class. Default is java.lang.String.

The Smile broker allows to reference topics simply by the topic name, which allows replacing a route like the following:

    <route>
        <from uri="kafka:v2-in-topic?brokers=localhost:9092&amp;sslKeystoreLocation=/path/to/keystore.jks&amp;sslKeystorePassword=changeit&amp;sslKeyPassword=changeit&amp;securityProtocol=SSL" />
        <to uri="smile:hl7v2/hl7v2ToFhirProcessor" />
        <to uri="kafka:bundle-out-topic?brokers=localhost:9092&amp;sslKeystoreLocation=/path/to/keystore.jks&amp;sslKeystorePassword=changeit&amp;sslKeyPassword=changeit&amp;securityProtocol=SSL" />
    </route>

by the simpler definition:

<route>
	<from uri="smile:clustermgr/broker?topic=v2-in-topic" />
	<to uri="smile:persistence/bundleProcessor" />
	<to uri="smile:clustermgr/broker?topic=bundle-out-topic" />
</route>

Example: Using Message Type and Payload Type Parameters

The following example shows how to use the messageType and payloadType parameters to handle typed messages. This is particularly useful when working with transaction log messages or other structured message types:

<route id="test-transaction-log">
	<from uri="smile:clustermgr/broker?topic=crd.events&amp;messageType=ca.cdr.broker.transaction.TransactionLogMessage&amp;payloadType=ca.cdr.api.model.json.TransactionLogEventsJson$TransactionLogEventJson"/>
	<log logName="ca.cdr.camel" message="${body}" loggingLevel="INFO"/>
	<to uri="bean:customProcessor"/>
</route>

In this example:

  • The messageType parameter specifies that we expect TransactionLogMessage objects
  • The payloadType parameter specifies the structure of the payload within those messages
  • This ensures type safety and proper message handling throughout the route

Message Keys for Ordering

When message ordering is important (e.g., processing HL7v2 feeds or related FHIR updates), you need to ensure messages about the same entity are routed to the same partition. This is done by setting a partition key (Kafka) or message key (Pulsar).

Using smile:clustermgr/broker

When using the smile:clustermgr/broker endpoint, the message key is determined by the IMessage.getMessageKey() method of the message being sent. To control the message key, you can either:

  1. Use a message type that implements custom key logic
  2. Set the key in a processor before sending

Using Apache Kafka Camel Component

When using the Apache Kafka Camel component directly (e.g., kafka:topic-name), you can set the partition key using Camel message headers:

<route>
    <from uri="direct:hl7v2-input"/>
    <process ref="patientIdExtractor"/>
    <to uri="kafka:hl7v2-topic?brokers=localhost:9092"/>
</route>

The processor extracts the patient identifier and sets it as the Kafka key:

import org.apache.camel.Exchange;
import org.apache.camel.Processor;
import org.apache.camel.component.kafka.KafkaConstants;
import ca.uhn.hl7v2.model.v25.datatype.CX;
import ca.uhn.hl7v2.model.v25.segment.PID;

public class PatientIdExtractorProcessor implements Processor {
    @Override
    public void process(Exchange exchange) throws Exception {
        org.apache.camel.Message message = exchange.getIn();
        ca.uhn.hl7v2.model.Message hapiMsg =
            message.getBody(ca.uhn.hl7v2.model.Message.class);

        // Extract patient ID from PID segment
        // Use both PID-3.1 (ID Number) and PID-3.4 (Assigning Authority)
        // to create a unique key across different identifier systems
        PID pid = (PID) hapiMsg.get("PID");
        CX patientIdentifier = pid.getPatientIdentifierList(0);
        String idNumber = patientIdentifier.getIDNumber().getValue();
        String assigningAuthority = patientIdentifier
            .getAssigningAuthority()
            .getNamespaceID().getValue();
        String partitionKey = assigningAuthority + "|" + idNumber;

        // Set as Kafka message key - ensures all messages for this
        // patient go to the same partition
        message.setHeader(KafkaConstants.KEY, partitionKey);
    }
}

Key Headers

BrokerHeaderDescription
KafkaKafkaConstants.KEYMessage key used for partition assignment
KafkaKafkaConstants.PARTITION_KEYExplicit partition number (use KEY instead for hash-based routing)
PulsarSet via message key in producerSee Pulsar documentation

For more information on message ordering concepts, see Message Ordering and Partition Keys.

Kafka Manual Commit

 

When your route begins with a Kafka consumer as the source, using manual commit mode can be useful in order to guarantee that no messages will be lost in the case of a disruption.

  • Example to URI: <to uri="smile:clustermgr/kafkaManualCommit" />

In order to use this processor, the Kafka consumer component must include the parameters autoCommitEnable=false and allowManualCommit=true.

The following example shows a route with manual transaction committing.

<route>
	<from uri="kafka:guaranteed-delivery-topic?brokers=localhost:9092&amp;allowManualCommit=true&amp;autoCommitEnable=false"/>
	<to uri="smile:persistence/bundleProcessor"/>
	<to uri="smile:clustermgr/kafkaManualCommit" />
</route>