【问题标题】:How can I send large messages with Kafka (over 15MB)?如何使用 Kafka(超过 15MB)发送大消息?
【发布时间】:2014-01-28 00:16:18
【问题描述】:

我使用 Java Producer API 将字符串消息发送到 Kafka V. 0.8。 如果消息大小约为 15 MB,我会收到 MessageSizeTooLargeException。 我尝试将 message.max.bytes 设置为 40 MB,但仍然出现异常。小消息没有问题。

(异常出现在生产者,我在这个应用程序中没有消费者。)

我能做些什么来摆脱这个异常?

我的示例生产者配置

private ProducerConfig kafkaConfig() {
    Properties props = new Properties();
    props.put("metadata.broker.list", BROKERS);
    props.put("serializer.class", "kafka.serializer.StringEncoder");
    props.put("request.required.acks", "1");
    props.put("message.max.bytes", "" + 1024 * 1024 * 40);
    return new ProducerConfig(props);
}

错误日志:

4709 [main] WARN  kafka.producer.async.DefaultEventHandler  - Produce request with correlation id 214 failed due to [datasift,0]: kafka.common.MessageSizeTooLargeException
4869 [main] WARN  kafka.producer.async.DefaultEventHandler  - Produce request with    correlation id 217 failed due to [datasift,0]: kafka.common.MessageSizeTooLargeException
5035 [main] WARN  kafka.producer.async.DefaultEventHandler  - Produce request with   correlation id 220 failed due to [datasift,0]: kafka.common.MessageSizeTooLargeException
5198 [main] WARN  kafka.producer.async.DefaultEventHandler  - Produce request with correlation id 223 failed due to [datasift,0]: kafka.common.MessageSizeTooLargeException
5305 [main] ERROR kafka.producer.async.DefaultEventHandler  - Failed to send requests for topics datasift with correlation ids in [213,224]

kafka.common.FailedToSendMessageException: Failed to send messages after 3 tries.
at kafka.producer.async.DefaultEventHandler.handle(Unknown Source)
at kafka.producer.Producer.send(Unknown Source)
at kafka.javaapi.producer.Producer.send(Unknown Source)

【问题讨论】:

  • 我的第一反应是要求您将这条巨大的信息分成几个较小的信息:-/ 我的猜测是由于某种原因这是不可能的,但您可能仍想重新考虑它:巨大消息通常意味着某处存在设计缺陷,应该真正修复。
  • 谢谢,但这会使我的逻辑复杂得多。为什么将 Kafka 用于 15MB 左右的消息是一个的主意? 1 MB 是可以使用的最大邮件大小限制吗?我在 Kafka 文档中发现的消息大小限制并不多。
  • 这与 Kafka 或任何其他消息处理系统完全无关。我的理由是:如果您的 15MB 文件出现问题,那么事后清理这些烂摊子是非常昂贵的。这就是为什么我通常将大文件拆分为许多较小的作业(然后通常也可以并行执行)。
  • 您是否使用过任何压缩方式?能否请您分享更多细节,仅凭一个词很难猜出一些东西
  • 对于那些偶然发现这个问题,但使用librdkafka与Kafka沟通的人,另请参阅:stackoverflow.com/questions/60739858/…

标签: java apache-kafka


【解决方案1】:

您需要调整三个(或四个)属性:

  • 消费者端:fetch.message.max.bytes - 这将确定消费者可以获取的最大消息大小。
  • 代理方:replica.fetch.max.bytes - 这将允许代理中的副本在集群内发送消息并确保消息被正确复制。如果这太小,则消息将永远不会被复制,因此,消费者将永远看不到该消息,因为该消息将永远不会被提交(完全复制)。
  • 代理方:message.max.bytes - 这是代理可以从生产者处接收到的最大消息大小。
  • 代理端(每个主题):max.message.bytes - 这是代理允许附加到主题的最大消息大小。此大小在压缩前经过验证。 (默认为经纪人的message.max.bytes。)

我发现了第 2 点的困难之处 - 您不会收到来自 Kafka 的任何异常、消息或警告,因此在发送大消息时请务必考虑这一点。

【讨论】:

  • 好的,你和 user2720864 是正确的。我只在源代码中设置了message.max.bytes。但是我必须在Kafka服务器config/server.properties的配置中设置这些值。现在更大的消息也起作用了:)。
  • 将这些值设置得太高有什么已知的缺点吗?
  • 是的。在消费者方面,您为每个分区分配fetch.message.max.bytes 内存。这意味着如果你使用大量的fetch.message.max.bytes 结合大量的分区,它会消耗大量的内存。实际上,由于broker之间的复制过程也是一个专门的消费者,这也会消耗broker上的内存。
  • 注意,还有一个max.message.bytes 配置per-topic 可以低于经纪人的message.max.bytes
  • 根据官方文档,消费者方面的参数和关于代理之间复制的参数/.*fetch.*bytes/ 似乎不是硬限制:“这不是绝对最大值,如果 [.. .]大于此值,仍会返回记录批次,以确保可以进行进度。”
【解决方案2】:

Kafka 0.10new consumerlaughing_man's answer 相比需要进行细微更改:

  • 经纪人:没有变化,你还是需要增加属性message.max.bytesreplica.fetch.max.bytesmessage.max.bytes 必须等于或小于 (*) replica.fetch.max.bytes
  • 生产者:增加max.request.size 以发送更大的消息。
  • 消费者:增加max.partition.fetch.bytes以接收更大的消息。

(*) 阅读 cmets 以了解更多关于 message.max.bytesreplica.fetch.max.bytes

【讨论】:

  • 你知道为什么message.max.bytes需要小于replica.fetch.max.bytes吗?
  • "replica.fetch.max.bytes(默认值:1MB)- 代理可以复制的最大数据大小。这必须大于 message。 max.bytes,否则代理将接受消息但无法复制它们。导致潜在的数据丢失。”来源:handling-large-messages-kafka
  • 感谢您回复我的链接。这似乎也与Cloudera guide 的建议相呼应。然而,这两个都是错误的 - 请注意,它们没有提供任何技术原因来说明 为什么 replica.fetch.max.bytes 应该严格大于 message.max.bytes。我怀疑的 Confluent 员工 confirmed earlier today:这两个数量实际上可以相等。
  • 是否有关于message.max.bytes<replica.fetch.max.bytesmessage.max.bytes=replica.fetch.max.bytes @Kostas 的更新?
  • 是的,它们可以相等:mail-archive.com/users@kafka.apache.org/msg25494.html(Ismael 为 Confluent 工作)
【解决方案3】:

@laughing_man 的回答非常准确。但是,我还是想给出一个我从 Kafka 专家Stephane Maarek 那里学到的建议。我们在我们的实时系统中积极应用了这个解决方案。

Kafka 不适合处理大型消息。

您的 API 应该使用云存储(例如,AWS S3),并简单地将对 S3 的引用推送到 Kafka 或任何其他消息代理。您需要找到一个地方来保存您的数据,无论它可以是网络驱动器还是完全其他的东西,但它不应该是消息代理。

如果您不想继续使用上述推荐的可靠解决方案,

消息最大大小为 1MB(您的代理中的设置称为message.max.bytesApache Kafka。如果您真的非常需要它,您可以增加该大小并确保为您的生产者和消费者增加网络缓冲区。

如果您真的关心拆分消息,请确保每个消息拆分具有完全相同的键,以便将其推送到同一分区,并且您的消息内容应报告“部分 id”,以便您的消费者可以完全重构消息。

如果消息是基于文本的,请尝试压缩数据,这可能会减少数据大小,但不会神奇。

同样,您必须使用外部系统来存储该数据,并且只需将外部引用推送到 Kafka。这是一种非常常见的架构,您应该采用并被广泛接受。

请记住,只有当消息量很大但不是很大时,Kafka 才能发挥最佳效果。

来源:https://www.quora.com/How-do-I-send-Large-messages-80-MB-in-Kafka

【讨论】:

  • Kafka 可以处理大消息,绝对没有问题。 Kafka 主页上的介绍页面甚至将其称为存储系统。
  • @Bhanu Hoysala - 我应该将大消息保存到存储中,然后在消息中发送参考。话虽如此,你如何保证数据被写入和引用消息被原子推送?要么都成功,要么都不成功。
  • @Jeremy 我们需要另一个主题/队列列出对存储桶所做的更改(我们可以配置为仅获取创建事件的通知)。在成功的情况下,我们将根据配置获取消息(您不会收到来自 S3 中失败操作的事件通知)。在失败的情况下,文件上传服务会知道写入是否成功(这是一个同步操作)。 docs.aws.amazon.com/AmazonS3/latest/dev/NotificationHowTo.html 取决于代理和存储组合,可以进行各种集成。
  • @Player_Neo,你说“Kafka 不是用来处理大消息的。”。您能否也阐明增加消息大小的影响?
【解决方案4】:

这个想法是让相同大小的消息从 Kafka Producer 发送到 Kafka Broker,然后由 Kafka Consumer 接收,即

Kafka 生产者 --> Kafka Broker --> Kafka 消费者

假设如果要求发送 15MB 的消息,那么 ProducerBrokerConsumer 三者都需要保持同步。

Kafka Producer 发送 15 MB --> Kafka Broker 允许/存储 15 MB --> Kafka Consumer 收到 15 MB

因此设置应该是:

a) 在代理上:

message.max.bytes=15728640 
replica.fetch.max.bytes=15728640

b) 关于消费者:

fetch.message.max.bytes=15728640

【讨论】:

  • 会不会是 ConsumerConfig 中的 fetch.message.max.bytes 被 max.partition.fetch.bytes 替代了?
  • 那你就不用换生产者了吧?
【解决方案5】:

您需要覆盖以下属性:

代理配置($KAFKA_HOME/config/server.properties)

  • replica.fetch.max.bytes
  • message.max.bytes

Consumer Configs($KAFKA_HOME/config/consumer.properties)
这一步对我不起作用。我将它添加到消费者应用程序中,它运行良好

  • fetch.message.max.bytes

重启服务器。

查看此文档以获取更多信息: http://kafka.apache.org/08/configuration.html

【讨论】:

  • 对于命令行使用者,我需要使用 --fetch-size= 标志。它似乎没有读取 consumer.properties 文件 (kafka 0.8.1) 。我还建议使用 compression.codec 选项从生产者端打开压缩。
  • Ziggy 的评论对我 kafka 0.8.1.1 有效。谢谢!
  • 会不会是 ConsumerConfig 中的 fetch.message.max.bytes 被 max.partition.fetch.bytes 替代了?
【解决方案6】:

要记住的一个关键点是message.max.bytes 属性必须与消费者的fetch.message.max.bytes 属性同步。获取大小必须至少与最大消息大小一样大,否则可能会出现生产者发送的消息大于消费者可以消费/获取的消息的情况。可能值得一看。
您使用的是哪个版本的 Kafka?还提供您获得的更多详细信息跟踪。日志中是否有类似 ...payload size of xxxx larger than 1000000 的内容?

【讨论】:

  • 我已经用更多信息更新了我的问题:Kafka 版本 2.8.0-0.8.0;现在我只需要制作人。
【解决方案7】:

我认为,这里的大多数答案都已经过时或不完整。

为了参考answer of Sacha Vetter(带有Kafka 0.10的更新),我想提供一些额外的信息和官方文档的链接。


生产者配置:

代理/主题配置:

  • message.max.bytes (Link) 可以设置,如果想增加代理级别的消息大小。但是,从文档中:“可以使用主题级别 max.message.bytes 配置为每个主题设置。”
  • max.message.bytes (Link) 可能会增加,如果只有一个主题应该能够接受较大的文件。不得更改代理配置。

我总是更喜欢主题受限的配置,因为我可以自己配置主题作为 Kafka 集群的客户端(例如,使用 admin client)。我可能对代理配置本身没有任何影响。


在上面的答案中,必要时提到了一些更多配置:

来自文档:“这不是一个绝对最大值,如果第一个非空分区中的第一个记录批大于这个值,仍然会返回记录批以确保进度可以制作。”

来自文档:“记录是由消费者分批获取的。如果获取的第一个非空分区中的第一个记录批大于此限制,该批仍将返回以确保消费者可以取得进步。”

来自文档:“记录是由消费者分批获取的,如果获取的第一个非空分区中的第一个记录批大于这个值,那么记录批仍然会返回给确保消费者能够取得进步。”


结论:关于获取消息的配置不需要更改处理消息,大于这些配置的默认值(在小型设置中进行了测试)。可能,消费者可能总是得到大小为 1 的批次。但是,必须设置第一个块中的两个配置,如前面的答案中所述。

此说明不应说明任何有关性能的内容,也不应作为设置或不设置这些配置的建议。必须根据具体的计划吞吐量和数据结构单独评估最佳值。

【讨论】:

    【解决方案8】:

    对于使用landoop kafka的人: 您可以在环境变量中传递配置值,例如:

    docker run -d --rm -p 2181:2181 -p 3030:3030 -p 8081-8083:8081-8083  -p 9581-9585:9581-9585 -p 9092:9092
     -e KAFKA_TOPIC_MAX_MESSAGE_BYTES=15728640 -e KAFKA_REPLICA_FETCH_MAX_BYTES=15728640  landoop/fast-data-dev:latest `
    

    如果您使用的是 rdkafka,则在生产者配置中传递 message.max.bytes,如下所示:

      const producer = new Kafka.Producer({
            'metadata.broker.list': 'localhost:9092',
            'message.max.bytes': '15728640',
            'dr_cb': true
        });
    

    同样,对于消费者而言,

      const kafkaConf = {
       "group.id": "librd-test",
       "fetch.message.max.bytes":"15728640",
       ... .. }                                                                                                                                                                                                                                                      
    

    【讨论】:

      【解决方案9】:

      这是我如何使用kafka-python==2.0.2成功发送高达 100mb 的数据:

      经纪人:

      consumer = KafkaConsumer(
          ...
          max_partition_fetch_bytes=max_bytes,
          fetch_max_bytes=max_bytes,         
      )
      

      生产者(见最后的最终解决方案):

      producer = KafkaProducer(
          ...
          max_request_size=KafkaSettings.MAX_BYTES,
      )
      

      然后:

      producer.send(topic, value=data).get()
      

      这样发送数据后,出现如下异常:

      MessageSizeTooLargeError: The message is n bytes when serialized which is larger than the total memory buffer you have configured with the buffer_memory configuration.

      最后我增加了buffer_memory(默认32mb)来接收另一端的消息。

      producer = KafkaProducer(
          ...
          max_request_size=KafkaSettings.MAX_BYTES,
          buffer_memory=KafkaSettings.MAX_BYTES * 3,
      )
      

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2019-01-16
        • 2020-05-19
        • 2022-10-24
        • 1970-01-01
        • 1970-01-01
        • 2019-05-30
        相关资源
        最近更新 更多