【问题标题】:spring kafka template producer performancespring kafka 模板生产者表现
【发布时间】:2018-04-28 01:46:25
【问题描述】:

我正在使用 Spring Kafka 模板来生成消息。而且它产生消息的速度太慢了。生成 15000 条消息大约需要 8 分钟。

以下是我创建 Kafka 模板的方式:

 @Bean
  public ProducerFactory<String, GenericRecord> highSpeedAvroProducerFactory(
      @Qualifier("highSpeedProducerProperties") KafkaProperties properties) {
    final Map<String, Object> kafkaPropertiesMap = properties.getKafkaPropertiesMap();
    System.out.println(kafkaPropertiesMap);
    kafkaPropertiesMap.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
    kafkaPropertiesMap.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, AvroGenericSerializer.class);
    return new DefaultKafkaProducerFactory<>(kafkaPropertiesMap);
  }

  @Bean
  public KafkaTemplate<String, GenericRecord> highSpeedAvroKafkaTemplate(
      @Qualifier("highSpeedAvroProducerFactory") ProducerFactory<String, GenericRecord> highSpeedAvroProducerFactory) {
    return new KafkaTemplate<>(highSpeedAvroProducerFactory);
  }

这是我使用模板发送消息的方式:

@Async("servicingPlatformUpdateExecutor")
  public void afterWrite(List<? extends Account> items) {
    LOGGER.info("Batch start:{}",items.size());
    for (Test test : items) {
        if (test.isOmega()) {

          ObjectKeyRecord objectKeyRecord = ObjectKeyRecord.newBuilder().setType("test").setId(test.getId()).build();
          LOGGER.info("build start, {}",test.getId());

          GenericRecord message = MessageUtils.buildEventRecord(
              schemaService.findSchema(topicName)
                  .orElseThrow(() -> new OmegaException("SchemaNotFoundException", topicName)), objectKeyRecord, test);
          LOGGER.info("build end, {}",account.getId());
          LOGGER.info("send Started , {}",account.getId());
          ListenableFuture<SendResult<String, GenericRecord>> future = highSpeedAvroKafkaTemplate.send(topicName, objectKeyRecord.toString(), message);
          LOGGER.info("send Done , {}",test.getId());
          future.addCallback(new KafkaProducerFutureCallback(kafkaSender, topicName, objectKeyRecord.toString(), message));
        }
    }
    LOGGER.info("Batch end}");

  }

生产者属性:

metric.reporters = []
metadata.max.age.ms = 300000
reconnect.backoff.ms = 50
sasl.kerberos.ticket.renew.window.factor = 0.8
bootstrap.servers = [***VALID BROKERS****))]
ssl.keystore.type = JKS
sasl.mechanism = GSSAPI
max.block.ms = 9223372036854775807
interceptor.classes = null
ssl.truststore.password = null
client.id = producer-1
ssl.endpoint.identification.algorithm = null
request.timeout.ms = 30000
acks = all
receive.buffer.bytes = 32768
ssl.truststore.type = JKS
retries = 2147483647
ssl.truststore.location = null
ssl.keystore.password = null
send.buffer.bytes = 131072
compression.type = none
metadata.fetch.timeout.ms = 60000
retry.backoff.ms = 100
sasl.kerberos.kinit.cmd = /usr/bin/kinit
buffer.memory = 800000000
timeout.ms = 30000
key.serializer = class org.apache.kafka.common.serialization.StringSerializer
sasl.kerberos.service.name = kafka
sasl.kerberos.ticket.renew.jitter = 0.05
ssl.trustmanager.algorithm = PKIX
block.on.buffer.full = false
ssl.key.password = null
sasl.kerberos.min.time.before.relogin = 60000
connections.max.idle.ms = 540000
max.in.flight.requests.per.connection = 10
metrics.num.samples = 2
ssl.protocol = TLS
ssl.provider = null
ssl.enabled.protocols = [TLSv1.2]
batch.size = 40000000
ssl.keystore.location = null
ssl.cipher.suites = null
security.protocol = SASL_SSL
max.request.size = 1048576
value.serializer = class com.message.serialization.AvroGenericSerializer
ssl.keymanager.algorithm = SunX509
metrics.sample.window.ms = 30000
partitioner.class = class org.apache.kafka.clients.producer.internals.DefaultPartitioner
linger.ms = 2

这是显示调用 kakfatemplate 发送方法需要几毫秒的日志:

2018-04-27 05:29:05.691 INFO  - testservice -  - UpdateExecutor-1 - com.test.testservice.adapter.batch.testsyncjob.UpdateWriteListener:70 - build start, 1
2018-04-27 05:29:05.691 INFO  - testservice -  - UpdateExecutor-1 - com.test.testservice.adapter.batch.testsyncjob.UpdateWriteListener:75 - build end, 1
2018-04-27 05:29:05.691 INFO  - testservice -  - UpdateExecutor-1 - com.test.testservice.adapter.batch.testsyncjob.UpdateWriteListener:76 - send Started , 1
2018-04-27 05:29:05.778 INFO  - testservice -  - UpdateExecutor-1 - com.test.testservice.adapter.batch.testsyncjob.UpdateWriteListener:79 - send Done , 1
2018-04-27 05:29:07.794 INFO  - testservice -  - kafka-producer-network-thread | producer-1 - com.test.testservice.adapter.batch.testsyncjob.KafkaProducerFutureCallback:38

任何关于如何提高发件人性能的建议将不胜感激

Spring Kakfa 版本:1.2.3.RELEASE 卡夫卡客户端:0.10.2.1

更新1:

将 Serializer 更改为 ByteArraySerializer 然后生成相同的。 我仍然看到对 kafkatempate 的每个发送方法调用需要 100 到 200 毫秒

ObjectKeyRecord objectKeyRecord = ObjectKeyRecord.newBuilder().setType("test").setId(test.getId()).build();
          GenericRecord message = MessageUtils.buildEventRecord(
              schemaService.findSchema(testConversionTopicName)
                  .orElseThrow(() -> new TestException("SchemaNotFoundException", testTopicName)), objectKeyRecord, test);
          byte[] messageBytes = serializer.serialize(testConversionTopicName,message);
          LOGGER.info("send Started , {}",test.getId());
          ListenableFuture<SendResult<String, byte[]>> future = highSpeedAvroKafkaTemplate.send(testConversionTopicName, objectKeyRecord.toString(), messageBytes);
          LOGGER.info("send Done , {}",test.getId());
          future.addCallback(new KafkaProducerFutureCallback(kafkaSender, testConversionTopicName, objectKeyRecord.toString(), message));

【问题讨论】:

    标签: apache-kafka kafka-producer-api spring-kafka


    【解决方案1】:

    您是否分析过您的应用程序?例如使用 YourKit。

    我怀疑它是 Avro 序列化程序;我能够在 274 毫秒内发送 15,000 条 1000 字节的消息。

    @SpringBootApplication
    public class So50060086Application {
    
        public static void main(String[] args) {
            SpringApplication.run(So50060086Application.class, args);
        }
    
        @Bean
        public ApplicationRunner runner(KafkaTemplate<String, String> template) {
            return args -> {
                Thread.sleep(5_000);
                String payload = new String(new byte[999]);
                StopWatch watch = new StopWatch();
                watch.start();
                for (int i = 0; i < 15_000; i++) {
                    template.send("so50060086a", "" + i + payload);
                }
                watch.stop();
                System.out.println(watch.prettyPrint());
            };
        }
    
        @Bean
        public NewTopic topic() {
            return new NewTopic("so50060086a", 1, (short) 1);
        }
    }
    

    StopWatch '': running time (millis) = 274
    

    【讨论】:

    • Gary ,我尝试更改为 ByteArraySerializer ,但仍然看到发送调用需要更多时间。用详细信息更新了我的帖子。尝试找到一些分析器,以便我可以使用它来查看发生了什么
    • 我刚刚注意到你有acks = all - 我刚刚更改了我的测试以使用它,我的测试达到了 436 毫秒(主题在 3 个代理上复制)。但是,我的 3 个经纪人在本地主机上。 acks=all 很贵——你有多少个代理实例?他们之间有良好的网络吗?
    • @Garry 我们有五个代理,每个主题在 3 个代理中复制。它们都在 AWS 同一个区域运行。
    • 可能是网络不佳的问题 - 无论是代理之间还是您的客户与代理之间。即使我在每次发送时将get() 添加到未来(等待所有确认),每次发送也只需要 1-2 毫秒。
    • 我还在思考为什么调用异步方法的 send 需要很长时间
    猜你喜欢
    • 2021-04-23
    • 1970-01-01
    • 2018-11-16
    • 2019-05-09
    • 1970-01-01
    • 1970-01-01
    • 2019-11-30
    • 1970-01-01
    • 2021-03-17
    相关资源
    最近更新 更多