Интеграция с Kafka
Описание
В данном разделе описывается пример работы с Kafka в Entaxy ION: настройка коннекторов, создание маршрутов и проверка обмена данными.
Процесс взаимодействия с Kafka состоит из следующих шагов:
-
Устанавливаем фичу camel-kafka:
feature:install -r camel-kafka -
Проверяем, что установленные бандлы находятся в статусе 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 -
Запускаем Kafka в докере:
docker pull apache/kafka:latest docker run -d --name=kafka -p 9092:9092 apache/kafka:latest -
Создаем профиль системы
kafkaи добавляем к нему:-
Входной кастомный коннектор, который будет читать сообщения из 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> -
Выходной кастомный коннектор, который отправляет сообщения в 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> -
К выходному кастомному коннектору добавляем маршрут на базе компонента таймер, инициирующий вызов коннектора с заданным интервалом:
<?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>
-
-
В логах отображается результат работы:
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