【发布时间】:2018-04-11 21:06:10
【问题描述】:
刚开始使用 Spring Kafka (2.1.4.RELEASE) 和 Kafka (1.0.0) 但是当我添加事务时,处理速度降低了很多。
代码:
spring.kafka.consumer.max-poll-records=10
spring.kafka.consumer.specific.avro.reader=true
spring.kafka.consumer.auto-offset-reset=earliest
spring.kafka.consumer.group-id=${application.name}
spring.kafka.consumer.properties.isolation.level=read_committed
spring.kafka.consumer.key-deserializer=io.confluent.kafka.serializers.KafkaAvroDeserializer
spring.kafka.consumer.value-deserializer=io.confluent.kafka.serializers.KafkaAvroDeserializer
我在 Java 中添加了:
@Bean
ProducerFactory<Object, Object> producerFactory(KafkaProperties properties) {
DefaultKafkaProducerFactory<Object, Object> factory = new DefaultKafkaProducerFactory<>(properties.buildProducerProperties());
factory.setTransactionIdPrefix(properties.getProducer().getTransactionIdPrefix());
return factory;
}
@Bean
KafkaTemplate<Object, Object> kafkaTemplate(ProducerFactory<Object, Object> factory) {
return new KafkaTemplate<>(factory, true);
}
@Bean("kafkaListenerContainerFactory")
ConcurrentKafkaListenerContainerFactory<Object, Object> listenerContainerFactory(Environment env, ConsumerFactory<Object, Object> consumerFactory, KafkaTransactionManager<Object, Object> transactionManager) {
ConcurrentKafkaListenerContainerFactory<Object, Object> factory = new ConcurrentKafkaListenerContainerFactory<>();
factory.setAutoStartup(true);
factory.setConcurrency(1);
factory.setConsumerFactory(consumerFactory);
factory.getContainerProperties().setTransactionManager(transactionManager);
factory.getContainerProperties().setGroupId(env.getRequiredProperty("spring.kafka.consumer.group-id"));
return factory;
}
当我删除setTransactionManager(transactionManager) 语句后,速度大幅提升。是不是我做错了什么?
【问题讨论】:
标签: performance apache-kafka kafka-consumer-api spring-kafka