Camel Kafka Connector

Try it out on Docker

This example runs Kafka in KRaft mode, builds a Kafka Connect image that includes the camel-http-sink connector package, and forwards messages from a Kafka topic to an HTTP endpoint.

It uses Camel Kafka Connector 4.18.0 (the latest package published to Maven Central at the time of writing) and Apache Kafka 3.9.1, which is the Kafka version that release depends on. Swap CKC_VERSION in the Dockerfile for a newer connector package from the Connectors List when you need one.

You need Docker and Docker Compose.

Prepare the example files

Create an empty directory and add the three files below.

Dockerfile
FROM apache/kafka:3.9.1

ARG CKC_VERSION=4.18.0
ARG CONNECTOR=camel-http-sink-kafka-connector
ARG BASEURL=https://repo1.maven.org/maven2/org/apache/camel/kafkaconnector

USER root
RUN mkdir -p /opt/kafka/plugins \
    && wget -qO- "${BASEURL}/${CONNECTOR}/${CKC_VERSION}/${CONNECTOR}-${CKC_VERSION}-package.tar.gz" \
       | tar -C /opt/kafka/plugins -xz \
    && chown -R appuser:appuser /opt/kafka/plugins
COPY connect-distributed.properties /opt/kafka/config/connect-distributed.properties
USER appuser

The connector archive unpacks to /opt/kafka/plugins/camel-http-sink-kafka-connector/, which is the layout Kafka Connect expects under plugin.path.

connect-distributed.properties
bootstrap.servers=kafka:19092
group.id=ckc-docker-example
key.converter=org.apache.kafka.connect.storage.StringConverter
value.converter=org.apache.kafka.connect.storage.StringConverter
offset.storage.topic=connect-offsets
offset.storage.replication.factor=1
config.storage.topic=connect-configs
config.storage.replication.factor=1
status.storage.topic=connect-status
status.storage.replication.factor=1
offset.flush.interval.ms=10000
plugin.path=/opt/kafka/plugins
listeners=HTTP://0.0.0.0:8083
rest.advertised.host.name=localhost
rest.advertised.port=8083
docker-compose.yml
services:
  kafka:
    image: apache/kafka:3.9.1
    hostname: kafka
    ports:
      - "9092:9092"
    environment:
      KAFKA_NODE_ID: 1
      KAFKA_PROCESS_ROLES: broker,controller
      KAFKA_LISTENERS: CONTROLLER://:29093,PLAINTEXT_HOST://:9092,PLAINTEXT://:19092
      KAFKA_ADVERTISED_LISTENERS: PLAINTEXT_HOST://localhost:9092,PLAINTEXT://kafka:19092
      KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT,PLAINTEXT_HOST:PLAINTEXT
      KAFKA_CONTROLLER_QUORUM_VOTERS: 1@kafka:29093
      KAFKA_INTER_BROKER_LISTENER_NAME: PLAINTEXT
      KAFKA_CONTROLLER_LISTENER_NAMES: CONTROLLER
      CLUSTER_ID: 4L6g3nShT-eMCtK--X86sw
      KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
      KAFKA_GROUP_INITIAL_REBALANCE_DELAY_MS: 0
      KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 1
      KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 1
      KAFKA_LOG_DIRS: /tmp/kraft-combined-logs

  kafka-connect:
    build: .
    hostname: kafka-connect
    depends_on:
      - kafka
    ports:
      - "8083:8083"
    entrypoint: ["/opt/kafka/bin/connect-distributed.sh"]
    command: ["/opt/kafka/config/connect-distributed.properties"]

  echo-server:
    image: mendhak/http-https-echo:31
    environment:
      HTTP_PORT: 80

Start the stack

From the directory that contains those files:

docker compose up --build

Wait until Kafka Connect answers on port 8083:

curl http://localhost:8083/connector-plugins

The response should list org.apache.camel.kafkaconnector.httpsink.CamelHttpsinkSinkConnector.

Create the HTTP sink connector

This example uses Kafka Connect in distributed mode, so the connector is created through the REST API.

connector.json
{
  "name": "http-sink-connector",
  "config": {
    "connector.class": "org.apache.camel.kafkaconnector.httpsink.CamelHttpsinkSinkConnector",
    "tasks.max": "1",
    "topics": "topic1",
    "key.converter": "org.apache.kafka.connect.storage.StringConverter",
    "value.converter": "org.apache.kafka.connect.storage.StringConverter",
    "camel.kamelet.http-sink.url": "http://echo-server/test",
    "camel.kamelet.http-sink.method": "PUT"
  }
}
curl -sS -H "Accept: application/json" -H "Content-Type: application/json" \
  -X POST --data @connector.json http://localhost:8083/connectors

Check that the connector task is running:

curl http://localhost:8083/connectors/http-sink-connector/status

Produce a message

echo YOUR_MESSAGE | docker compose exec -T kafka \
  /opt/kafka/bin/kafka-console-producer.sh --bootstrap-server localhost:9092 --topic topic1

The echo service should log a PUT /test request whose body is YOUR_MESSAGE.

To use a different Camel Kafka Connector package, change CKC_VERSION (and CONNECTOR if you are not using HTTP) in the Dockerfile, rebuild, and set connector.class plus the camel.kamelet.* keys from that connector’s reference page.