【问题标题】:Test a Reactive-Kafka Consumer and Producer Template using embedded kafka + custom serialised使用嵌入式 kafka + 自定义序列化测试 Reactive-Kafka 消费者和生产者模板
【发布时间】:2021-09-07 06:55:40
【问题描述】:

我们需要一个关于如何使用embedded-kafka-broker 测试ReactiveKafkaConsumerTemplateReactiveKafkaProducerTemplate 的示例。谢谢。

正确的代码在这里讨论后

你可以有你的自定义de-serializer相应地使用自定义ReactiveKafkaConsumerTemplate

自定义序列化器:


import org.apache.kafka.common.serialization.Serializer;

import com.fasterxml.jackson.core.JsonProcessingException;
import com.fasterxml.jackson.databind.ObjectMapper;

public class EmployeeSerializer implements Serializer<Employee> {

    @Override
    public byte[] serialize(String topic, Employee data) {
        
        byte[] rb = null;
        ObjectMapper mapper = new ObjectMapper();
        try {
            rb = mapper.writeValueAsString(data).getBytes();
        } catch (JsonProcessingException e) {
            e.printStackTrace();
        }
        return rb;
    }

}

在 Embedded-kfka-reactive 测试中使用它:


import java.util.Map;

import org.apache.kafka.clients.producer.ProducerConfig;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.apache.kafka.connect.json.JsonSerializer;
import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.springframework.kafka.core.reactive.ReactiveKafkaProducerTemplate;
import org.springframework.kafka.support.converter.MessagingMessageConverter;
import org.springframework.kafka.test.condition.EmbeddedKafkaCondition;
import org.springframework.kafka.test.context.EmbeddedKafka;
import org.springframework.kafka.test.utils.KafkaTestUtils;

import reactor.kafka.sender.SenderOptions;
import reactor.kafka.sender.SenderRecord;
import reactor.test.StepVerifier;

@EmbeddedKafka(topics = EmbeddedKafkareactiveTest.REACTIVE_INT_KEY_TOPIC,
brokerProperties = { "transaction.state.log.replication.factor=1", "transaction.state.log.min.isr=1" })
public class EmbeddedKafkareactiveTest {

    public static final String REACTIVE_INT_KEY_TOPIC = "reactive_int_key_topic";

    private static final Integer DEFAULT_KEY = 1;

    private static final String DEFAULT_VERIFY_TIMEOUT = null;

    private ReactiveKafkaProducerTemplate<Integer, Employee> reactiveKafkaProducerTemplate;

    @BeforeEach
    public void setUp() {
        reactiveKafkaProducerTemplate = new ReactiveKafkaProducerTemplate<>(setupSenderOptionsWithDefaultTopic(),
                new MessagingMessageConverter());
    }

    private SenderOptions<Integer, Employee> setupSenderOptionsWithDefaultTopic() {
        Map<String, Object> senderProps = KafkaTestUtils
                .producerProps(EmbeddedKafkaCondition.getBroker().getBrokersAsString());
        SenderOptions<Integer, Employee> senderOptions = SenderOptions.create(senderProps);
        senderOptions = senderOptions.producerProperty(ProducerConfig.TRANSACTIONAL_ID_CONFIG, "reactive.transaction")
                .producerProperty(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, true)
                .producerProperty(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, JsonSerializer.class.getName())
                ;
        return senderOptions;
    }

    @Test
    public void test_When_Publish() {
        
        
        Employee employee = new Employee();
        
        ProducerRecord<Integer, Employee> producerRecord = new ProducerRecord<Integer, Employee>(REACTIVE_INT_KEY_TOPIC, DEFAULT_KEY, employee);
                
        StepVerifier.create(reactiveKafkaProducerTemplate.send(producerRecord)
                .then())
                .expectComplete()
                .verify();
    }   

    @AfterEach
    public void tearDown() {
        reactiveKafkaProducerTemplate.close();
    }
}

【问题讨论】:

    标签: spring-kafka reactive-kafka


    【解决方案1】:

    添加了使用非事务性生产者进行的正确序列化。请查看本页顶部的代码以获取答案。

    【讨论】:

      【解决方案2】:

      框架中的测试使用嵌入式 kafka 代理。

      https://github.com/spring-projects/spring-kafka/tree/main/spring-kafka/src/test/java/org/springframework/kafka/core/reactive

      @EmbeddedKafka(topics = ReactiveKafkaProducerTemplateIntegrationTests.REACTIVE_INT_KEY_TOPIC, partitions = 2)
      public class ReactiveKafkaProducerTemplateIntegrationTests {
      ...
      

      【讨论】:

      • 不要把代码放在cmets中;最好编辑您所做的问题和评论。我需要查看完整的测试和实际的异常。
      • 谢谢。我已经添加了。
      • 请查看我的编辑以获取正确的降价以格式化代码。 setupSenderOptionsWithDefaultTopic() - 如果你从框架测试中复制了它,那么,是的,它为键设置了一个 IntegerSerializer,为值设置了 SrtringSerializer - 你需要一个可以处理你的 Employee 对象的值序列化程序,例如JsonSerializer.
      • “不起作用”没有传达任何有用的信息。如果你的意思是你仍然得到同样的错误,这意味着你没有正确设置属性。如果设置正确,它将起作用。我无法从代码 sn-p 中看出你在做什么。添加完整的测试,而不仅仅是 sn-ps。
      • 这不是同一个问题,这是一个不同的错误。有两个问题。 1.您使用了错误的JsonSerializer;它应该是import org.springframework.kafka.support.serializer.JsonSerializer;。 2. 您正在创建一个事务性生产者并使用非事务性发送。您的测试通过了非事务性生产者和正确的序列化程序。
      猜你喜欢
      • 2019-05-09
      • 1970-01-01
      • 1970-01-01
      • 2014-04-09
      • 2018-01-07
      • 2020-05-21
      • 2019-01-15
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多