【发布时间】: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