【问题标题】:Problems with Kafka configuration using SpringBoot使用 Spring Boot 配置 Kafka 的问题
【发布时间】:2018-12-26 16:11:19
【问题描述】:

我们的团队一直在试验 Kafka 的问题。自从我们开始开发我们的应用程序以来,这个问题就一直存在。

一开始,这些问题很快就解决了。 “它失败了?只需重新启动服务器”。由于我们想开始向公众分发我们的应用程序,因此这种“解决方案”不再可行。

我们面临的问题基本上有两个:

消费者停止工作

这是一个经常性的。突然之间,一些消费者就停下来了。消息成功发送到Kafka,我们甚至可以使用Kafka Tool看到实际的消息,但是consumer不起作用。

循环消息

这是相反的。有时会发送一条消息,而消费者只是继续消费该消息,直到我们重新启动服务器。

我们尝试将 Kafka 直接配置到服务器中,但我们意识到由于某种原因 Kafka 会忽略这些配置并直接从 Spring Boot 中获取配置。

我们的配置如下所示:

消费者:

    properties.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, KAFKA_BOOTSTRAP);
    properties.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
    properties.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
    properties.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false);
    properties.put(ConsumerConfig.GROUP_ID_CONFIG, KAFKA_CONSUMER_GROUP);
    properties.put(ConsumerConfig.MAX_PARTITION_FETCH_BYTES_CONFIG, 5000000);
    properties.put(ConsumerConfig.REQUEST_TIMEOUT_MS_CONFIG, 10800000);

制作人:

    properties.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, this.KAFKA_HOST);
    properties.put(ProducerConfig.RETRIES_CONFIG, 0);
    properties.put(ProducerConfig.BATCH_SIZE_CONFIG, 16384);
    properties.put(ProducerConfig.LINGER_MS_CONFIG, 1);
    properties.put(ProducerConfig.BUFFER_MEMORY_CONFIG, 33554432);
    properties.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
    properties.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, JsonSerializer.class);
    properties.put(ProducerConfig.MAX_REQUEST_SIZE_CONFIG, 5000000);
    properties.put(ProducerConfig.COMPRESSION_TYPE_CONFIG, "gzip");
    properties.put(ProducerConfig.REQUEST_TIMEOUT_MS_CONFIG, 10800000);

您可以看到极端的超时和大小值,这是因为我们认为问题与消息的大小或服务器的超时有关。我们甚至重新设计了我们所有的应用程序流量,以便我们可以发送更小的消息,那时我们意识到这不是根本问题。

任何帮助将不胜感激。

卡夫卡 0.10 版

Spring Boot 版本 1.5.7 Dalston.SR4

我们正在使用 spring-cloud-starter-stream-kafka 依赖项。

我们检查了日志,但没有可识别的错误。实际上,只有信息消息,但没有一个能说明有用的信息。

【问题讨论】:

  • 什么版本的 Boot 和 spring-kafka?当前版本分别为 2.0.4 和 2.1.7。什么版本的卡夫卡经纪人?您需要显示更多配置 - 即您使用的是 Boot 的自动配置工厂还是您自己的工厂等。您是否进行了线程转储以查看容器线程在做什么?日志中是否有任何内容表明存在问题?
  • 我编辑了这篇文章。谢谢。

标签: java spring-boot apache-kafka spring-kafka


【解决方案1】:

不幸的是,Boot 1.5 引入了一个不再受支持的非常旧版本的 spring-kafka (1.1.x)。 Boot 对依赖版本控制有严格的规定。如the Spring for Apache Kafka project page 所述:

推荐所有 brokers >= 0.10.x.x 的用户使用 spring-kafka 1.3.x 或更高版本,因为它的线程模型更简单,这要归功于 KIP-62。

当前的 1.3.x 版本是1.3.5。尝试升级到该版本和 0.11 kafka-clients jar。

在 1.1.x 中,有复杂的逻辑需要缓慢的侦听器来暂停/恢复消费者以避免代理重新平衡。虽然我自己没有看到,但我看到了消费者在暂停后无法正常恢复的报告。

感谢 KIP-62,不再需要这种逻辑,因为心跳是在后台发送的。但是,您需要确保max.poll.interval.ms 足够大以支持处理max.poll.records

如果可以的话,升级到 boot 2.0.4 会更好,它会引入最新的 spring-kafka 2.1.7。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2019-07-14
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多