Ответ
Работал с Apache Kafka в рамках тестирования event-driven микросервисных архитектур. Основные активности:
- Интеграционное и E2E-тестирование: Проверка корректности отправки, маршрутизации и обработки событий между сервисами.
- Валидация данных: Контроль схемы сообщений (использовал Apache Avro с Schema Registry), порядка следования (ordering) и семантики доставки (at-least-once, exactly-once).
- Тестирование отказоустойчивости: Проверка поведения при падении консьюмеров, обработка сообщений в Dead Letter Queue (DLQ), ретраи.
- Мониторинг и отладка: Использовал Kafkacat, UI for Apache Kafka и Grafana с метриками Kafka для наблюдения за лагом консьюмеров и состоянием топиков.
Пример автотеста на Java (Kafka Clients API) для проверки отправки и получения:
// Настройка продюсера
Properties producerProps = new Properties();
producerProps.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
producerProps.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
producerProps.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
producerProps.put(ProducerConfig.ACKS_CONFIG, "all"); // Гарантия доставки
try (Producer<String, String> producer = new KafkaProducer<>(producerProps)) {
// Отправка синхронно с проверкой
RecordMetadata metadata = producer.send(
new ProducerRecord<>("orders", "order-123", "{"status":"new"}")
).get();
System.out.println("Sent to partition: " + metadata.partition());
}
// Настройка консьюмера для валидации
Properties consumerProps = new Properties();
consumerProps.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
consumerProps.put(ConsumerConfig.GROUP_ID_CONFIG, "test-group");
consumerProps.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
consumerProps.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
consumerProps.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
try (Consumer<String, String> consumer = new KafkaConsumer<>(consumerProps)) {
consumer.subscribe(List.of("orders"));
ConsumerRecords<String, String> records = consumer.poll(Duration.ofSeconds(10));
// Проверка полученного сообщения
records.forEach(record -> {
assert record.key().equals("order-123");
assert record.value().contains(""status":"new"");
});
}