Быстрый старт#

Предварительные требования:

  • 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)