【问题标题】:Understanding Kafka Message Byte Size了解 Kafka 消息字节大小
【发布时间】:2019-11-02 15:59:41
【问题描述】:

如何获取 Kafka 中单条记录的大小?

有一些关于我为什么需要这个的说明。

这似乎不是 ConsumerRecord 或 RecordMetadata 类上公开的 serializedValueSize。我不太了解这个属性的价值,因为它与对消费者有用的消息的大小不匹配。如果不是这个,serializedValueSize 是做什么用的?

我试图让我的 Kafka java 应用程序表现得像“min.poll.records”,如果它存在以补充“max.poll.records”。我必须这样做,因为它是必需的:)。假设给定主题上的所有消息都具有相同的大小(在这种情况下是正确的),这应该可以从消费者方面通过将 fetch.min.bytes 设置为消息量乘以每个消息的字节大小来实现。消息。

这是存在的:

https://kafka.apache.org/documentation/#consumerapi

max.poll.records

在一次 poll() 调用中返回的最大记录数。

这不存在,但这是我想要的行为:

min.poll.records

在一次 poll() 调用中返回的最小记录数。如果在 fetch.max.wait.ms 中指定的时间过去之前没有足够的记录可用,则无论如何都会返回记录,因此这不是绝对最小值。

这是我目前发现的:

  • 在生产者方面,我将“batch.size”设置为 1 个字节。这会强制生产者单独发送每条消息。

  • 关于消费者大小,我将“max.partition.fetch.bytes”设置为 291 字节。这使得消费者只能得到 1 条消息。将此值设置为 292 会使消费者有时会收到 2 条消息。所以我计算出消息大小是 292 的一半; 一条消息的大小为 146 字节

  • 上述项目符号需要更改 Kafka 配置并涉及手动查看 / grepping 一些服务器日志。如果 Kafka Java API 提供了这个值,那就太好了。

  • 在生产者方面,Kafka 提供了一种方法来获取RecordMetadata.serializedValueSize method 中记录的序列化大小。这个值是 76 字节,与上面测试中给出的 146 字节有很大不同。

  • 在消费者规模上,Kafka 提供了ConsumerRecord API。此记录的序列化值大小也是 76。偏移量每次仅增加 1(而不是记录的字节大小)。

  • key的大小为-1字节(key为null)。

System.out.println(myRecordMetadata.serializedValueSize());
// 76
# producer
batch.size=1

# consumer

# Expected this to work:
# 76 * 2 = 152
max.partition.fetch.bytes=152

# Actually works:
# 292 = ??? magic ???
max.partition.fetch.bytes=292

我希望将 max.partition.fetch.bytes 设置为 serializedValueSize 给出的字节数的倍数将使 Kafka 消费者从轮询中接收最多该数量的记录。相反,max.partition.fetch.bytes 值需要更高才能发生这种情况。

【问题讨论】:

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


    【解决方案1】:

    原答案

    我对@9​​87654325@ 方法不太熟悉,但根据文档,这只是该消息中存储的值的大小。这将小于总消息大小(即使使用 null 键),因为消息还包含不属于值的元数据(例如时间戳)。

    至于您的问题:与其直接通过处理消息大小和限制消费者的吞吐量来控制轮询,为什么不只是缓冲传入的消息,直到有足够的可用消息或所需的超时(您提到fetch.max.wait.ms 但是你可以手动指定一个)已经过去了吗?

    public static <K, V> List<ConsumerRecord<K, V>>
        minPoll(KafkaConsumer<K, V> consumer, Duration timeout, int minRecords) {
      List<ConsumerRecord<K, V>> acc = new ArrayList<>();
      long pollTimeout = Duration.ofMillis(timeout.toMillis()/10);
      long start = System.nanoTime();
      do {
        ConsumerRecords<K, V> records = consumer.poll(pollTimeout);
        for(ConsumerRecord<K, V> record : records)
          acc.add(record);
      } while(acc.size() < minRecords &&
              System.nanoTime() - start < timeout.toNanos());
      return acc;
    }
    

    consumer.poll 的调用中的timeout.toMillis()/10 超时是任意的。您应该选择一个足够小的持续时间,这样即使我们等待的时间长于指定的超时时间(这里:长 10%)也没关系。

    编辑:请注意,这可能会返回一个大于max.poll.records 的列表(最大值为max.poll.records + minRecords - 1)。如果您还需要强制执行此严格的上限,请使用该方法外部的另一个缓冲区来临时存储多余的记录(这可能会更快,但不允许混合使用 minPoll 和普通的 poll 方法)或只需丢弃它们并使用consumerseek 方法回溯即可。

    回答更新的问题

    所以问题不在于控制poll-方法返回的消息数量,而在于如何获取单个记录的大小。不幸的是,我认为如果不经历很多麻烦,这是不可能的。问题是这个问题没有真正的(恒定的)答案,甚至一个大概的答案也取决于 Kafka 版本,或者更确切地说是不同的 Kafka 协议版本。

    首先,我不完全确定max.partition.fetch.bytes 究竟控制了什么(例如:协议开销是否也是其中的一部分?)。让我说明一下我的意思:当消费者发送一个 fetch 请求时,那么 fetch 响应由以下字段组成:

    1. 限制时间(4 个字节)
    2. 主题响应数组(4 个字节的数组长度 + 数组中的数据大小)。

    主题响应依次包括

    1. 主题名称(字符串长度 + 字符串大小 2 个字节)
    2. 分区响应数组(4 个字节的数组长度 + 数组中的数据大小)。

    然后有一个分区响应

    1. 分区 ID(4 个字节)
    2. 错误代码(2 个字节)
    3. 高水印(8 字节)
    4. 最后的稳定偏移量(8 个字节)
    5. 日志起始偏移量(8 个字节)
    6. 中止事务的数组(数组长度为 4 个字节 + 数组中的数据)
    7. 记录集。

    所有这些都可以在FetchResponse.java 文件中找到。一个记录集又由包含记录的记录批组成。我不会列出包含记录批次的所有内容(您可以看到它here)。只需说开销为 61 个字节即可。最后,批处理中单个记录的大小有点棘手,因为它使用 varint 和 varlong 字段。它包含

    1. 正文大小(1-5 个字节)
    2. 属性(1 字节)
    3. 时间戳增量(1-10 字节)
    4. 偏移增量(1-5 个字节)
    5. 关键字节数组(1-5字节+关键数据大小)
    6. 值字节数组(1-5字节+值数据大小)
    7. 标头(1-5 字节 + 标头数据大小)。

    源代码是here。如您所见,您不能简单地将 292 字节除以 2 来获得记录大小,因为某些开销是恒定的并且与记录数无关。

    更糟糕的是记录没有固定大小,即使它们的键和值(和标头)有,因为时间戳和偏移量存储为与批处理时间戳和偏移量的差异,使用可变长度数据类型。此外,这只是撰写本文时最新协议版本的情况。对于旧版本,答案将再次不同,谁知道未来版本会发生什么。

    【讨论】:

    • 谢谢。我对此表示赞同,因为这是一种明智的方法。我没有单击接受,因为它与我要问的问题太不同了,即“如何在 Kafka 中获取单个记录的大小?”。我编辑了主要问题以反映这是我的主要目标。
    • 你使用的是哪个 Kafka 版本?
    • 使用 Kafka 2.2.0。
    猜你喜欢
    • 2015-03-06
    • 2020-12-04
    • 2018-10-19
    • 1970-01-01
    • 2021-02-18
    • 1970-01-01
    • 2015-10-26
    • 2018-10-16
    • 1970-01-01
    相关资源
    最近更新 更多