我创建了以下测试类来使用 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
我创建了一个单元测试来测试 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)