小编Sma*_*hie的帖子

SpringBoot 嵌入式 Kafka 使用 Avro Schema 生成事件

我创建了以下测试类来使用 AvroSerializer 生成事件。

@SpringBootTest
@EmbeddedKafka(partitions = 1, brokerProperties = { "listeners=PLAINTEXT://localhost:9092", "port=9092" })
@TestPropertySource(locations = ("classpath:application-test.properties"))
@ContextConfiguration(classes = { TestAppConfig.class })
@DirtiesContext
@TestInstance(TestInstance.Lifecycle.PER_CLASS)
class EntitlementEventsConsumerServiceImplTest {


    @Autowired
    EmbeddedKafkaBroker embeddedKafkaBroker;

    @Bean
    MockSchemaRegistryClient mockSchemaRegistryClient() {
        return new MockSchemaRegistryClient();
    }

    @Bean
    KafkaAvroSerializer kafkaAvroSerializer() {
        return new KafkaAvroSerializer(mockSchemaRegistryClient());
    }

    @Bean
    public DefaultKafkaProducerFactory producerFactory() {
        Map<String, Object> props = KafkaTestUtils.producerProps(embeddedKafkaBroker);
        props.put(KafkaAvroSerializerConfig.AUTO_REGISTER_SCHEMAS, false);
        return new DefaultKafkaProducerFactory(props, new StringSerializer(), kafkaAvroSerializer());
    }

    @Bean
    public KafkaTemplate<String, ApplicationEvent> kafkaTemplate() {
        KafkaTemplate<String, ApplicationEvent> kafkaTemplate = new KafkaTemplate(producerFactory());
        return kafkaTemplate;
    }
} …
Run Code Online (Sandbox Code Playgroud)

java apache-kafka spring-kafka confluent-schema-registry spring-kafka-test

6
推荐指数
2
解决办法
4564
查看次数

嵌入式 Kafka 测试 - 未调用 Kafka 监听器

我创建了一个单元测试来测试 Kafka 监听器,如下所示。

@SpringBootTest
@EmbeddedKafka(partitions = 1, brokerProperties = {"listeners=PLAINTEXT://localhost:9092","port=909"})
class ConsumerTest {

    @Autowired
    KafkaTemplate producer;

    @Test
    public void consumeEvents1Test() throws InterruptedException {
        producer.send("events1", "Sample message");
        Thread.sleep(1000);
    }
}
Run Code Online (Sandbox Code Playgroud)

Consumer 的创建如下所示。

@Component
public class Consumer {
    Logger LOG = LoggerFactory.getLogger(Consumer.class);

    @KafkaListener(id= "${topic1}" ,
            topics = "${topic1}",
            groupId = "${consumer.group1}", concurrency = "1", containerFactory = "kafkaListenerContainerFactory")
    public void consumeEvents1(String message, @Headers Map<String, String> header, Acknowledgment acknowledgment) {
        LOG.info("Message - {}", message);
        LOG.info(header.get(KafkaHeaders.GROUP_ID) + header.get(KafkaHeaders.RECEIVED_TOPIC)+String.valueOf(header.get(KafkaHeaders.OFFSET)));
        acknowledgment.acknowledge();


    }
}
Run Code Online (Sandbox Code Playgroud)

消费者工厂和容器工厂的创建如下所示。

@Bean
public ConsumerFactory<String, …
Run Code Online (Sandbox Code Playgroud)

apache-kafka spring-boot spring-kafka spring-kafka-test

3
推荐指数
1
解决办法
5418
查看次数