【问题标题】:Spring + Kafka: Transactions slowSpring + Kafka:事务缓慢
【发布时间】: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


    【解决方案1】:

    Kafka 事务非常昂贵 - 特别是如果您提交每个发送。

    Transactions in Apache Kafka

    向下滚动到“事务如何执行,以及如何调整它们”。

    正如我们所见,开销与作为事务的一部分写入的消息数量无关。因此,获得更高吞吐量的关键是在每个事务中包含更多的消息。

    使用 Spring for Apache Kafka,您可以使用 executeInTransaction 方法在同一事务中进行多次发送。或者通过 KafkaTransactionManager 使用 Spring 事务管理并在 @Transactional 方法中执行多次发送。

    编辑

    我没有注意到监听器容器;我假设您正在消费一条消息,执行一些转换并发送到另一个主题。因此,在这种情况下,您不能“在事务中发送多条消息”,因为容器管理事务,并且默认情况下,在每次交付后提交。

    增加并发不会影响事务语义;在您的情况下,(并发 10),分区分布在 10 个线程中。每个线程运行一个单独的事务。

    您可以通过在容器工厂上将batchListener 设置为true 来进一步加快速度。

    在这种情况下,您的@KafkaListener 将获得List&lt;ConsumerRecord&gt;(或List&lt;Foo&gt;,如果您正在使用转换);您可以遍历列表并处理每条记录并使用模板发送(不要使用executeInTransaction,因为已经有一个事务,由容器线程启动)。然后,当批处理完成时,容器将提交事务。

    您可以使用 kafka max.poll.records consuer 属性控制批量大小。

    【讨论】:

    • 我明白了。我尝试将其更改为批处理模式,一次轮询 100 条记录和并发 10 条记录。这使它更快,但是事务将如何工作呢?是一次还是?如果它在例如之后失败5条记录,交易包裹在哪里? 100 左右读取,10 并发还是?
    • 查看我的答案的编辑;我没有注意到您使用的是交易型消费者;在那种情况下,情况有点不同。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2018-05-01
    • 2021-10-31
    • 1970-01-01
    • 1970-01-01
    • 2019-11-30
    • 2020-03-07
    相关资源
    最近更新 更多