Починаємо роботу з Apache Kafka. Частина IV
Всім привіт. Я Сергій Моренець, розробник, викладач, спікер та технічний письменник, хочу поділитися з вами своїм досвідом роботи з такою цікавою технологією, як Apache Kafka та розкрити ті теми, які з нею пов’язані.
У попередніх статтях я розповів про внутрішні особливості Kafka, конфігурацію та деплоймент відправників повідомлень та тестування Kafka додатків.
Це дуже цікава тема, яку ми розглядаємо на деяких тренінгах. Наприклад, візьмемо ту саму абстрактну систему, побудовану на мікросервісах, у якій платіжний сервіс керує платежами. І кожен платіж (успішний чи невдалий) повинен призводити до відправки події (нотифікації), для чого ми використовуємо Kafka.
Ці нотифікації будуть оброблятися у двох інших мікросервісах:
- квитковий — оформляє квиток;
- сервіс поїздок — оновлює кількість вільних місць на рейсі.
Spring Kafka
Як і раніше, можна використовувати Kafka клієнт для Java, але оскільки ці наші сервіси працюють на основі Spring Boot, то правильніше було б використовувати інтеграцію Spring та Apache Kafka — Spring Kafka. Цей проєкт, який додає ще один рівень абстракції над Kafka-клієнтом, дозволяє використовувати Kafka у знайомому нам стилі, з бінами, автоконфігурацією, властивостями, які можна прописувати в properties/YAML файлах та багато іншого.
У Spring Kafka є два підходи для написання конфігурації:
- Декларативний (YAML).
- Імперативний, через явні оголошення Spring бінів та їх налаштування.
Другий підхід більш гнучкий і єдиний можливий, якщо ви не використовуєте Spring Boot. Але оскільки у нас є Spring Boot проєкт, то ми можемо використовувати декларативну конфігурацію. Додамо необхідні налаштування у файл application.properties для сервісу поїздок:
spring.kafka.consumer.bootstrap-servers[0]=kafka:9092
spring.kafka.consumer.group-id=trip
spring.kafka.consumer.auto-offset-reset=earliest
spring.kafka.consumer.key-deserializer=org.apache.kafka.common.serialization.StringDeserializer
spring.kafka.consumer.value-deserializer=org.springframework.kafka.support.serializer.JsonDeserializer
Нагадаю, що для відправлення повідомлень у Kafka обов’язково була лише одна властивість — boostrap.servers. Але для отримання повідомлень кількість обов’язкових властивостей набагато більша. Розберемо їх по одному:
1)bootstrap.servers — адреса Kafka брокера (або брокерів), до якої потрібно буде підключатися.2) group.id — унікальний ідентифікатор споживача у концепції Kafka. Справа в тому, що в Kafka повідомлення можуть читати різні споживачі та ідентифікуються вони якраз group.id. Якби у нас був не один розділ, а три, наприклад, то ми могли б підключити три екземпляри одного і того ж сервісу, які б паралельно обробляли повідомлення, але у них би group.id був би один. У термінології Kafka вони належали б до однієї consumer group.
3) auto.offset.reset — це важлива властивість, яка визначає, звідки новому споживачеві читати повідомлення: від початку (earliest) або з останнього повідомлення (latest). Щоб споживачі не читали те саме повідомлення двічі, Kafka зберігає для них поточний offset (зміщення останнього прочитаного повідомлення) у спеціальному topic. І таким чином Kafka знає, яке повідомлення споживач прочитав останнім. Для нового споживача ця інформація відсутня, тому Kafka використовує цю властивість.
4) key.deserializer — клас, який використовується для десеріалізації ключа повідомлення.
5) value.deserializer- клас (з Spring Kafka), який використовується для десеріалізації значення повідомлення. Як ви бачите, тут такий параметр доводиться вказувати вручну.
Тепер додамо Java конфігурацію:
@EnableKafka
public class KafkaConfig {}
Spring Kafka вимагає анотації @EnableKafka, щоб увімкнути автоконфігурацію для тих бінів, які відносяться до Kafka. Як ми отримуватимемо повідомлення? У Spring є зручна концепція слухачів (listeners), які отримують Spring події з application context. Аналогічно в Spring Kafka можна отримувати повідомлення в спеціальному класі-обробнику, причому ви самі вибираєте, ви будете їх отримувати по одному або в групі (batch). Для початку виберемо найпростіший варіант, коли ми отримуємо та обробляємо одне повідомлення:
public class KafkaEventConsumer {@KafkaListener(topics = "payments")
public void consume(final ConsumerRecord<?, ? extends BaseEvent<?>> record) {}
}
Анотація @KafkaListener дозволяє підписатися на будь-яке повідомлення з певного topic’a (або topic’ів). При цьому саме повідомлення буде обернене до спеціального класу ConsumerRecord з Kafka клієнта. Цей клас цікавий тим, що його можна параметризувати (типами ключа та значення) і він містить всю інформацію про повідомлення:
- Ключ.
- Значення.
- Номер розділу (partition).
- Час відправлення.
- Зміщення (offset) повідомлення у topic.
- Назва topic’a.
При цьому нам доведеться самим десеріалізувати значення повідомлення з байтового масиву шуканий Java об’єкт. Такий підхід актуальний тоді, коли ми обробляємо повідомлення довільного вмісту. Але в даному випадку ми хочемо отримати лише один тип події — Payment Success. І це можна зробити, замінивши аргумент ConsumerRecord на PaymentSuccessEvent і додавши анотацію @Payload, щоб Spring зрозумів, що цей аргумент потрібно брати зі значення повідомлення:
@KafkaListener(topics = "payments")
public void consumePaymentSuccessEvent(final @Payload PaymentSuccessEvent event) {
У цій реалізації є один мінус. Якщо під час обробки повідомлення трапиться програмна помилка, повідомлення залишиться прочитаним, але не обробленим. Вдруге прочитати його не вийде, тому що Kafka змістить поточний offset для нашого споживача. Вихід із цієї ситуації є — додати спеціальний аргумент Acknowledgement, який дозволить вручну підтверджувати прочитання повідомлення за допомогою методу acknowledge:
public void consumePaymentSuccessEvent(final @Payload PaymentSuccessEvent event,
Acknowledgment acknowledgment) { try {acknowledgment.acknowledge();
} catch (Exception e) {log.error(e.getMessage(), e);
}
Тут головне не забути встановити властивість ackMode в manual, тому що інакше Spring Kafka думатиме, що потрібно автоматично позначати повідомлення прочитаними відразу після отримання:
spring.kafka.listener.ack-mode=manual
Але як дізнатися ключ повідомлення (у разі, якщо це ідентифікатор замовлення)? За допомогою спеціальної анотації @Header, яка дозволить отримати певні метадані з повідомлення (у даному випадку ключ):
public void consumePaymentSuccessEvent(final @Payload PaymentSuccessEvent event,
@Header(name = KafkaHeaders.RECEIVED_MESSAGE_KEY) Integer key,
Acknowledgment acknowledgment) {
Якщо зараз запустити наш додаток у Docker контейнері та відправити платіжне повідомлення, то у лозі Kafka контейнера буде дуже багато схожих помилок:
CreateTopics result(s): CreatableTopic(name=’__consumer_offsets’, numPartitions=50, replicationFactor=3, assignments=[], configs=[CreateableTopicConfig(name=’compression.type’, value=’producer’), CreateableTopicConfig(name=’cleanup.policy’, value=’compact’), CreateableTopicConfig(name=’segment.bytes’, value=’104857600′)]): INVALID_REPLICATION_FACTOR (Unable to replicate the partition 3 time(s): The target replication factor of 3 cannot be reached because only 1 broker(s) are registered.) (org.apache.kafka.controller.ReplicationControlManager)
Детальний аналіз показує, що коли сервіс поїздок підключається до Kafka брокеру, щоб прочитати повідомлення, а topic’a ще немає, він автоматично створюється. Але проблема в тому, що в налаштуваннях за замовчуванням (у Kafka) зазначено replication factor 3:
offsets.topic.replication.factor = 3
offsets.topic.num.partitions = 50
А така конфігурація можлива лише якщо у нас три Kafka брокери (а не один, як зараз). З іншого боку, за замовчуванням вибирається 50 розділів (partitions) на topic, що дуже багато для нашого навчального варіанту, де кількість передплатників буде вкрай невелика. А кожен розділ забирає дисковий простір і знижує ефективність. Можна сміливо змінити це значення на 3 у docker-compose.yml:
kafka:
build:
context: docker-scripts/kafka
environment:
KAFKA_NODE_ID: 1
KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: 'CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT'
KAFKA_ADVERTISED_LISTENERS: 'PLAINTEXT://kafka:9092'
KAFKA_PROCESS_ROLES: 'broker,controller'
KAFKA_CONTROLLER_QUORUM_VOTERS: '1@kafka:9093'
KAFKA_LISTENERS: 'PLAINTEXT://kafka:9092,CONTROLLER://kafka:9093'
KAFKA_CONTROLLER_LISTENER_NAMES: 'CONTROLLER'
KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
KAFKA_OFFSETS_TOPIC_NUM_PARTITIONS: 3
Перезбираємо Docker образи, запускаємо контейнери. І отримуємо вже нову помилку:
Caused by: org.apache.kafka.common.errors.RecordDeserializationException: Error deserializing key/value for partition payments-0 at offset 0. If needed, please seek past the record to continue consumption.
at org.apache.kafka.clients.consumer.internals.Fetcher.parseRecord(Fetcher.java:1429) ~[kafka-clients-3.0.1.jar!/:na]
at org.apache.kafka.clients.consumer.internals.Fetcher.access$3400(Fetcher.java:134) ~[kafka-clients-3.0.1.jar!/:na]
... 3 common frames omitted
Caused by: java.lang.IllegalStateException: No type information in headers and no default type provided
at org.springframework.util.Assert.state(Assert.java:76) ~[spring-core-5.3.19.jar!/:5.3.19]
at org.springframework.kafka.support.serializer.JsonDeserializer.deserialize(JsonDeserializer.java:583) ~[spring-kafka-2.8.5.jar!/:2.8.5]
Kafka і type mappings
Щоб зрозуміти причину цієї помилки, ще раз розглянемо код, який викликається при отриманні нового повідомлення:
public class KafkaEventConsumer {@KafkaListener(topics = "payments")
public void consumePaymentSuccessEvent(final @Payload PaymentSuccessEvent event, @Header(name = KafkaHeaders.RECEIVED_MESSAGE_KEY) Integer key,
Acknowledgment acknowledgment) {
Тут тип аргументу PaymentSuccessEvent, а це означає, що Spring Kafka повинен автоматично десеріалізувати значення отриманого повідомлення, конвертувати його з байтового масиву в Java-об’єкт. Але лише в тому випадку, якщо це повідомлення зберігає цей об’єкт. Звідси виникає питання. Як це дізнається Spring Kafka? Адже в самому повідомленні «тип» ніде не зберігається.
На жаль, у Kafka немає стандартного механізму опису метаданих про повідомлення. Ми часто використовуємо RESTful веб-сервіси, в яких завдяки HTTP протоколу є можливість додавати заголовки запиту/відповіді. І коли ми додаємо новий REST-сервіс:
@PostMapping(path = "orders", consumes = MediaType.APPLICATION_JSON_VALUE)
public void create(@RequestBody @Valid CreateOrderDTO createOrderDTO) {то тут можна вказати, який формат (JSON) тіла запиту ми очікуємо чи надсилаємо клієнту. Більше того, тут немає проблеми ідентифікації типу, тому що ми його вказуємо як аргумент методу (CreateOrderDTO), і шлях «orders» і HTTP метод POST однозначно трактує те, що цей метод повинен прийти тільки той запит, у якого тип CreateOrderDTO. Будь-який інший тип (а це легко з’ясується при десеріалізації) буде помилкою.
У Kafka це зробити складніше, тому що навіть в один topic можуть надходити повідомлення різного типу. Спочатку Kafka не мала заголовків, але їх додали у версії 0.11 як можливість додатків описувати метадані про повідомлення. Тобто сама програма відповідає як формат значення в повідомленні, так і за його тип.
У Spring Kafka вчинили такий спосіб. Тут є спеціальна властивість JsonSerializer.ADD_TYPE_INFO_HEADERS, яка за замовчанням дорівнює true. Коли ми відправляємо повідомлення брокеру, Spring автоматично додає в метадані заголовок з назвою __TypeId__ і значенням — повний шлях до класу. Коли Spring Kafka обробляє таке повідомлення, воно автоматичне із заголовків бере назву класу та успішно десеріалізує. Але у нас повідомлення надсилаються у платіжному сервісі, за допомогою Micronaut Kafka, який нічого не знає про цю конвенцію та про ці заголовки.
Тому найбільш правильне і просте рішення — надсилати такий заголовок із платіжного сервісу. Але якщо ми просто відправлятимемо шлях класу, це буде не дуже гнучкий похід, тому що назва класу та його повний шлях можуть змінитися (наприклад, в результаті рефакторингу), що може поламати обробку. Тому краще відправляти спеціальний нейтральний «дискримінатор», який можна однозначно ототожнити з типом. І такий токен у нас є — це поле type у класі BaseEvent, яке ми ще ніяк не використовували.
@Getter
public abstract class BaseEvent<T> {private final String type;
Тому ми його відправлятимемо у вигляді такого заголовка, а завдання Spring Kafka прочитати цей заголовок і знайти в конфігурації шуканий клас.
Необхідно зробити ще одну зміну в сигнатурі методу-приймача:
@KafkaListener(topics = "payments")
public void consumePaymentSuccessEvent(final @Payload PaymentSuccessEvent event, @Header(name = KafkaHeaders.RECEIVED_MESSAGE_KEY) Integer key,
Acknowledgment acknowledgment) {Що буде, якщо відправити в топік orders повідомлення зі значенням типу PaymentFailureEvent? Ми зараз у сервісі поїздок таке повідомлення не обробляємо, але воно існує і обов’язково буде доставлене у topic orders, а значить, і прочитане всіма сервісами-підписниками.На жаль, у Spring Kafka немає механізму фільтрації повідомлень. І тип аргументу @Payload (PaymentSuccessEvent) він використовує тільки для десеріалізації. Тому єдиний варіант — це вказати базовий тип для всіх повідомлень (BaseEvent), а потім уже в run-time перевіряти отриманий тип (тим більше, що що ми вже маємо pattern matching для instanceof і pattern matching для switch):
public void consumePaymentSuccessEvent(final @Payload BaseEvent<?> event,
@Header(name = KafkaHeaders.RECEIVED_MESSAGE_KEY) Integer key, Acknowledgment acknowledgment) { try { switch(event) {case PaymentSuccessEvent success -> tripService.adjustAvailableSeats(key, 1);
Type mappings і десеріалізація
Наша система гетерогенна та використовує різні платформи (Spring Boot, Jakarta EE, Micronaut). І не дуже хочеться прив’язуватися до якоїсь стандартної назви, оскільки вона може відрізнятися для різних платформ. На жаль, додати заголовок для Kafka повідомлення з довільною назвою не так просто, тому що в Spring Kafka він hard-coded і змінити його — нетривіальне завдання, принаймні в нашому випадку:
public abstract class AbstractJavaTypeMapper implements BeanClassLoaderAware {public static final String DEFAULT_CLASSID_FIELD_NAME = "__TypeId__";
Тому ми у спрощеному варіанті виберемо саме __TypeId__ як назву заголовка.
public class MessagingUtil {private static final String HEADER_TYPE_ID = "__TypeId__";
public static String getTypeIdHeader() {return HEADER_TYPE_ID;
}
}
Наступне завдання — зробити так, щоб платіжний сервіс надсилав цей заголовок, причому відправляв завжди і автоматично, без необхідності програмісту втручатися в цей процес. Ми не можемо змінювати абстрактні методи в PaymentClient (і їх сигнатуру), але ми можемо додати дефолтний метод, який завжди додаватиме заголовок з типом повідомлення:
@KafkaClient(id="payment")
@Topic("payments")public interface PaymentClient {CompletableFuture<RecordMetadata> send(@KafkaKey Object key, BaseEvent<?> value, Headers headers);
}
default CompletableFuture<RecordMetadata> send(Object key, BaseEvent<?> value) {return send(key, value, List.of(
new RecordHeader(MessagingUtil.getTypeIdHeader(), value.getType().getBytes(StandardCharsets.UTF_8))));
}
Для цього в Micronaut є спеціальний API у вигляді інтерфейсу Headers, який, по суті, є Map, де ключ — це рядок, а значення — список (значень заголовків). Споживачі нашого PaymentClient можуть використовувати дефолтний метод send(), де заголовки вже встановлені, або абстрактний, і додавати їх самостійно.
Зараз значення типу повідомлення зберігатиметься у перерахуванні PaymentEventType:
public enum PaymentEventType {PAYMENT_SUCCESS, PAYMENT_FAILURE
}
Краще прийняти конвенцію, що тип — це завжди рядок у нижньому регістрі, де слова поділяються не символом підкреслення, а точкою. Відповідно, потрібно оновити конструктор класу BaseEvent:
public BaseEvent(String entityId, String type, String source, T payload) {this.entityId = entityId;
this.type = type.toLowerCase().replaceAll("_", ".");this.source = source;
this.payload = payload;
createdAt = LocalDateTime.now();
id = UUID.randomUUID().toString();
}
Тепер потрібно додати спеціальний type mapping до конфігураційного файлу application.properties:
spring.kafka.consumer.properties.spring.json.type.mapping=payment.success: payment.event.PaymentSuccessEvent,payment.failure: payment.event.PaymentFailureEvent
У Spring Kafka прийнято конвенцію, де така карта типів містить їх список, розділений комою, де спочатку йде назва типу, потім двокрапка і потім повний шлях до класу.
Висновки
Налаштування Kafka брокера не є таким банальним процесом, і вимагає хорошого знання архітектури Kafka. Крім того, при використанні Kafka у базовому варіанті необхідно вказувати в метаданих (заголовках) той тип значення, який відправляється в Kafka topic. Такі метадані підтримуються у Spring Kafka, але вимагають явної вказівки у Micronaut. Тим не менш, нам вдалося домогтися того, що повідомлення успішно відправляються та приймаються у гетерогенній системі.
2 коментарі
Додати коментар Підписатись на коментаріВідписатись від коментарівКласні статті, дякую!
Дякую за гарні статті!