Быстрый старт#
Предварительные требования:
JDK 11;
Дистрибутив Corax;
Любая IDE, позволяющая работать с Java-проектами. В примере используется IntelliJ Idea.
Подключение#
Подключите клиентскую библиотеку Corax, добавив в CLASSPATH проекта и подключив к проекту вместе с библиотекой slf4j:
kafka-clients-XXX.jar
slf4j-api-1.7.36.jar
Где XXX — версия дистрибутива Corax.
Использование#
В терминах Corax существуют следующие сущности:
Producer, поставщик данных;
Consumer, потребитель данных;
Topic, топик, в который пишут поставщики и из которого читают потребители;
Кластер Corax, ответственный за хранение, репликацию и т.д.
В примере создается поставщик данных, который будет писать данные о количестве машин-нарушителей, проехавших в данный момент мимо камеры наблюдения на улице Колотушкина. Поставщик пишет в топик offenders.
Нарушители появляются с некоторой вероятностью, за что будет отвечать генерация случайного числа. После того как было зафиксировано 10 нарушений – поставщик данных прекращает свою работу.
После этого создается потребитель данных, который будет собирать данные о нарушителях и писать их в консоль.
Считается, что кластер Corax уже создан и функционирует. Топик
offendersсоздан администратором.
Для простоты считается, что кластер поднят как PLAINTEXT, т.е без защиты ssl и т.д.
Producer (поставщик)#
Пример поставщика данных:
package ru.sbt.example;
import org.apache.kafka.clients.producer.KafkaProducer;
import org.apache.kafka.clients.producer.Producer;
import org.apache.kafka.clients.producer.ProducerConfig;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.apache.kafka.clients.producer.RecordMetadata;
import org.apache.kafka.common.serialization.IntegerSerializer;
import org.apache.kafka.common.serialization.StringSerializer;
import java.util.Properties;
import java.util.Random;
import java.util.concurrent.ExecutionException;
public class CarProducer {
public static void main(String[] args) throws ExecutionException, InterruptedException {
Random random = new Random();
String server = "localhost:9092";
String topicName = "offenders";
final Properties props = new Properties();
props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, server);
props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, IntegerSerializer.class);
final Producer<String, Integer> producer = new KafkaProducer<>(props);
int numberOfOffendersCaught = 0;
while (numberOfOffendersCaught < 10) {
if (random.nextInt(10) > 7) {
int offenders = random.nextInt(100);
System.out.println("Поймано на ул. Колотушкина: " + offenders);
RecordMetadata recordMetadata = producer.send(
new ProducerRecord<>(topicName, "колотушкина", offenders)).get();
if (recordMetadata.hasOffset())
System.out.println("Данные о нарушителях отправлены успешно");
numberOfOffendersCaught++;
} else {
System.out.println("На ул. Колотушкина пока никого...");
}
}
producer.close();
}
}
Где в переменной server, укажите адрес кластера (брокер).
В примере брокер поднят на том же самом хосте, где будет выполняться программа CarProducer, поэтому указан localhost.
Обратите внимание, что название улицы является ключом – это позволяет Corax группировать сообщения по партициям на брокере.
Запустите пример, должны появиться записи вида:
На ул. Колотушкина пока никого...
На ул. Колотушкина пока никого...
Поймано на ул. Колотушкина: 122
Данные о нарушителях отправлены успешно
На ул. Колотушкина пока никого...
Что говорит об успешной отправке данных.
Consumer (потребитель)#
Создайте потребителя, который бесконечно будет вычитывать данные из топика.
package ru.sbt.example;
import org.apache.kafka.clients.consumer.Consumer;
import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.clients.consumer.ConsumerRecords;
import org.apache.kafka.clients.consumer.KafkaConsumer;
import org.apache.kafka.common.serialization.IntegerDeserializer;
import org.apache.kafka.common.serialization.StringDeserializer;
import java.time.Duration;
import java.util.Collections;
import java.util.Properties;
public class CarConsumer {
public static void main(String[] args) {
String server = "localhost:9092";
String topicName = "offenders";
final Properties props = new Properties();
props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, server);
props.put(ConsumerConfig.GROUP_ID_CONFIG, "test-consumer-group");
props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, IntegerDeserializer.class);
props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
final Consumer<String, Integer> consumer = new KafkaConsumer<>(props);
consumer.subscribe(Collections.singletonList(topicName));
while (true) {
final ConsumerRecords<String, Integer> consumerRecords = consumer.poll(Duration.ofSeconds(5));
consumerRecords.forEach(record ->
System.out.printf("Кол-во нарушителей: (%s, %d)\n", record.key(), record.value())
);
consumer.commitAsync();
}
}
}
Запустите потребителя, данные из топика нарушителей должны выводиться на экран в формате:
Кол-во нарушителей: (колотушкина, 88)