Подписка на события#
Для передачи событий через компонент EVTD необходимо реализовать отправку событий и их получение на стороне клиентов: продюсер на стороне отправителя и консьюмер на стороне получателя.
Подписка на события – процесс получения сообщений из определенного топика определенного кластера Platform V Corax / Apache Kafka.
Для этого необходимо:
подготовить клиентский сертификат для подключения приложения к Platform V Corax / Apache Kafka;
определить кластер, куда будут публиковаться события, к которым необходимо получить доступ;
определить наименование топика, в который публикуются события, к которым необходимо получить доступ;
запросить у администратора доступа выдачу прав на подписку к указанному топику для полученного сертификата.
Шаги разработки:
Подключить в зависимости проекта библиотеку
"org.apache.kafka:kafka-clients:3.7.2"В классе произвести заполнение параметров подключения:
Properties props = new Properties();
props.put("bootstrap.servers", "{ IP_ADDRESS },{ IP_ADDRESS },{ IP_ADDRESS }");
props.put("group.id", "test");
props.put("enable.auto.commit", "true");
props.put("auto.commit.interval.ms", "1000");
props.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
props.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
props.put("security.protocol", "SSL");
props.put("ssl.keystore.location", "<полный путь до хранилища ключа>" );
props.put("ssl.truststore.location", "<полный путь до хранилища ключа>");
props.put("ssl.keystore.password", "<пароль хранилища ключа>");
props.put("ssl.truststore.password", "<пароль хранилища сертификата ЦС>");
props.put("ssl.enabled.protocols", "TLSv1.2");
props.put("ssl.key.password", "<пароль закрытого ключа>");
props.put("ssl.endpoint.identification.algorithm", "");
, где:
bootstrap.servers— список серверов кластера Platform V Corax / Apache Kafka с указанием портов через запятую;group.id— наименования группы потребителя, все потребители в рамках одной группы считаются одним потребителем;enable.auto.commit— установка вtrueозначает, что обработанные offsets коммитятся автоматически, интервал commit устанавливается в параметреauto.commit.interval.ms;key.deserializerиvalue.deserializer— класс-десериализатор ключа и значения, полученных вConsumerRecords. В комплекте с клиентом идут два стандартных классаByteArrayDeserializerиStringDeserializer.
Добавить подключение к кластеру:
KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props);
consumer.subscribe(Arrays.asList("my-topic", "your-topic"));
, где:
props— предзаполненные параметры подключения;my-topic,your-topic— наименование топиков, к которым производится подключение.
Добавить непосредственно вычитку:
while (true) {
ConsumerRecords<String, String> records = consumer.poll(100);
for (ConsumerRecord<String, String> record : records)
System.out.printf("offset = %d, key = %s, value = %s%n", record.offset(), record.key(), record.value());
}