【发布时间】:2020-06-23 15:37:59
【问题描述】:
我们正在使用 Cucumber 对 out 应用程序进行一些集成测试,并且在测试 @KafkaListener 时遇到了一些问题。我们设法使用 EmbeddedKafka 并在其中生成数据。
但消费者从未收到任何数据,我们也不知道发生了什么。
这是我们的代码:
生产者配置
@Configuration
@Profile("test")
public class KafkaTestProducerConfig {
private static final String SCHEMA_REGISTRY_URL = "schema.registry.url";
@Autowired
protected EmbeddedKafkaBroker embeddedKafka;
@Bean
public Map<String, Object> producerConfig() {
Map<String, Object> props = new HashMap<>();
props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG,
embeddedKafka.getBrokersAsString());
props.put(SCHEMA_REGISTRY_URL, "URL");
props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, KafkaAvroSerializer.class);
return props;
}
@Bean
public ProducerFactory<String, GenericRecord> producerFactory() {
return new DefaultKafkaProducerFactory<>(producerConfig());
}
@Bean
public KafkaTemplate<String, GenericRecord> kafkaTemplate() {
return new KafkaTemplate<>(producerFactory());
}
}
消费者配置
@Configuration
@Profile("test")
@EnableKafka
public class KafkaTestConsumerConfig {
@Autowired
protected EmbeddedKafkaBroker embeddedKafka;
private static final String SCHEMA_REGISTRY_URL = "schema.registry.url";
@Bean
public Map<String, Object> consumerProperties() {
Map<String, Object> props = new HashMap<>();
props.put(ConsumerConfig.GROUP_ID_CONFIG, "groupId");
props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, embeddedKafka.getBrokersAsString());
props.put(SCHEMA_REGISTRY_URL, "URL");
props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false);
props.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, "1000");
props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, KafkaAvroDeserializer.class.getName());
props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
props.put(ConsumerConfig.FETCH_MAX_WAIT_MS_CONFIG, 10000);
return props;
}
@Bean
public DefaultKafkaConsumerFactory<String, Object> consumerFactory() {
KafkaAvroDeserializer avroDeserializer = new KafkaAvroDeserializer();
avroDeserializer.configure(consumerProperties(), false);
return new DefaultKafkaConsumerFactory<>(consumerProperties(), new StringDeserializer(), avroDeserializer);
}
@Bean
public ConcurrentKafkaListenerContainerFactory<String, Object> kafkaListenerContainerFactory() {
ConcurrentKafkaListenerContainerFactory<String, Object> factory =
new ConcurrentKafkaListenerContainerFactory<>();
factory.setConsumerFactory(consumerFactory());
factory.setBatchListener(true);
factory.getContainerProperties().setAckMode(ContainerProperties.AckMode.MANUAL_IMMEDIATE);
return factory;
}
}
集成测试
@SpringBootTest(
webEnvironment = SpringBootTest.WebEnvironment.RANDOM_PORT,
classes = Application.class)
@ActiveProfiles("test")
@EmbeddedKafka(topics = {"TOPIC1", "TOPIC2", "TOPIC3"})
public class CommonStepDefinitions implements En {
protected static final Logger LOGGER = LoggerFactory.getLogger(CommonStepDefinitions.class);
@Autowired
protected KafkaTemplate<String, GenericRecord> kafkaTemplate;
}
步骤定义
public class KafkaStepDefinitions extends CommonStepDefinitions {
private static final String TEMPLATE_TOPIC = "TOPIC1";
public KafkaStepDefinitions(){
Given("given statement", () -> {
OperationEntity operationEntity = new OperationEntity();
operationEntity.setFoo("foo");
kafkaTemplate.send(TEMPLATE_TOPIC, AvroPojoTransformer.pojoToRecord(operationEntity));
});
}
}
消费者 同样的代码在生产 Bootstrap 服务器上运行良好,但嵌入式 Kafka 从未达到过
@KafkaListener(topics = "${kafka.topic1}", groupId = "groupId")
public void consume(List<GenericRecord> records, Acknowledgment ack) throws DDCException {
LOGGER.info("Batch of {} records received", records.size());
//do something with the data
ack.acknowledge();
}
日志中的一切看起来都很好,但我们不知道缺少什么。
提前致谢。
【问题讨论】:
-
对我来说一切都很好;你确定记录是公开的吗?你能把日志贴在某个地方(pastebin 等)吗?
-
日志中唯一看起来不好的是这个错误:删除时出错 ...\AppData\Local\Temp\kafka-7678436650062051819\TOPIC1-0\00000000000000000000.index:进程无法访问该文件,因为它正被另一个进程使用
-
是的 - 这是关闭期间 Windows 上的一个(尚未解决的)问题。如果您收到
partitions assignedINFO 日志,我几乎可以保证问题出在发布方面。再次,如果您可以发布日志,我可以看看。 -
这是集成测试的完整日志:pastebin.com/FjYmJ2TW
标签: java spring-boot apache-kafka spring-kafka spring-kafka-test