Интеграция с Kafka

Описание

В данном разделе описывается пример работы с Kafka в Entaxy ION: настройка коннекторов, создание маршрутов и проверка обмена данными.

Процесс взаимодействия с Kafka состоит из следующих шагов:

  1. Устанавливаем фичу camel-kafka:

    feature:install -r camel-kafka
  2. Проверяем, что установленные бандлы находятся в статусе Active:

    711 │ Active │ 50  │ 0        │ lz4-java
    712 │ Active │ 50  │ 3.4.5    │ camel-kafka
    713 │ Active │ 50  │ 2.4.1.1  │ Apache ServiceMix :: Bundles :: kafka-clients
    714 │ Active │ 120 │ 1.1.7.3  │ snappy-java: A fast compression/decompression library
  3. Запускаем Kafka в докере:

    docker pull apache/kafka:latest
    docker run -d --name=kafka -p 9092:9092 apache/kafka:latest
  4. Создаем профиль системы kafka и добавляем к нему:

    1. Входной кастомный коннектор, который будет читать сообщения из Kafka. В точке кастомизации custom-route создаем маршрут, который подписывается на топик test и выводит полученные сообщения из Kafka в лог:

      <?xml version="1.0" encoding="UTF-8"?>
      <entaxy:object-input-route
          xmlns="http://camel.apache.org/schema/blueprint"
          xmlns:blueprint="http://www.osgi.org/xmlns/blueprint/v1.0.0"
          xmlns:entaxy="http://www.entaxy.ru/schemas/1.0"
          xmlns:m="http://www.entaxy.ru/schemas/entaxy-mediators/1.0">
          <from uri="kafka:test?brokers=localhost:9092"/>
          <m:log message="In Kafka consumer : ${body}"/>
          <m:respond now="true" continue="false"/>
      </entaxy:object-input-route>
    2. Выходной кастомный коннектор, который отправляет сообщения в Kafka. В точке кастомизации custom-route создаем маршрут, который отправляет сообщения в топик test:

      <?xml version="1.0" encoding="UTF-8"?>
      <entaxy:object-route
          xmlns="http://camel.apache.org/schema/blueprint"
          xmlns:blueprint="http://www.osgi.org/xmlns/blueprint/v1.0.0"
          xmlns:entaxy="http://www.entaxy.ru/schemas/1.0"
          xmlns:m="http://www.entaxy.ru/schemas/entaxy-mediators/1.0">
          <m:log message="Custom Kafka output route start"/>
          <setBody>
              <constant>To Kafka</constant>
          </setBody>
          <m:log message="Kafka producer started: ${body}"/>
          <to uri="kafka:test?brokers=localhost:9092"/>
      </entaxy:object-route>
    3. К выходному кастомному коннектору добавляем маршрут на базе компонента таймер, инициирующий вызов коннектора с заданным интервалом:

      <?xml version="1.0" encoding="UTF-8"?>
      <entaxy:common-route
          xmlns="http://camel.apache.org/schema/blueprint"
          xmlns:blueprint="http://www.osgi.org/xmlns/blueprint/v1.0.0"
          xmlns:entaxy="http://www.entaxy.ru/schemas/1.0"
          xmlns:m="http://www.entaxy.ru/schemas/entaxy-mediators/1.0">
          <m:log message="in route ${routeId}" loggingLevel="INFO"/>
          <to uri="system:kafka?connectorClassifier=main"/>
      </entaxy:common-route>
  5. В логах отображается результат работы:

    18:03:57.749 INFO [Camel (kafka.custom-connector-out.main) thread #123 - timer://timer] in route kafka.custom-connector-out.main__timer
    18:03:57.750 INFO [Thread-994] Custom Kafka output route start 18:03:57.750 INFO [Thread-994] Kafka producer started: To Kafka
    18:03:57.762 INFO [Camel (kafka.custom-connector-in.main) thread #120 - KafkaConsumer[test]] c7313283-70a6-4e1b-8908-c909328bd89a#{"service":"sys-26","sender":"kafka"}# In Kafka consumer : To Kafka