Kafka Streams设置Header时Int/Long转Array[Byte]的正确方法
Kafka Streams写入Header时Int/Long转字节数组读取为空的解决方法
问题现象
基于Scala使用Kafka Streams给Topic记录添加Header时,按接口要求需将Int、Long类型值转换为Array[Byte]存储,但测试时两种常规数值转字节数组的方式写入后,通过kafkacat的%h参数读取Header时值为空,仅将数值转为UTF-8字符串再取字节数组的方式可正常读到值:
- 直接调用Int的
toByte方法构造字节数组写入counter_1,读取结果为空 - 通过
ByteBuffer.allocate(4).putInt(counter).array()转换后写入counter_2,读取结果为空 - 先将Int转为String,再取
StandardCharsets.UTF_8编码字节数组写入counter_3,可正常读到值1
kafkacat读取输出片段:
...counter_1=,counter_2=,counter_3=1,....
测试代码如下:
val counter: Int = 1 // counter_1 读取为空 context.headers().add("counter_1", Array[Byte](counter.toByte)) // counter_2 读取为空 context.headers().add("counter_2", ByteBuffer.allocate(4).putInt(counter).array()) // counter_3 可正常读取到值1 context.headers().add("counter_3", String.valueOf(counter).getBytes(StandardCharsets.UTF_8))
根本原因
前两种二进制转换方式的写入逻辑本身没有问题,Header值实际已经成功写入,读取为空是kafkacat的默认显示逻辑导致的:kafkacat打印Header时,默认会把字节数组按UTF-8编码解码为可打印字符串,不可打印的控制字符、空字符会被直接忽略不显示:
- 第一种
counter.toByte的写法:Int值1转为单个Byte后对应ASCII控制字符SOH(编码值0x01),属于不可打印控制字符,解码后无可见内容。另外这种写法本身有逻辑缺陷:仅取Int的最低8位,数值超过255时会发生溢出,得到错误结果,生产环境禁止使用。 - 第二种ByteBuffer转换的写法:Int值1的大端序4字节表示为
[0x00, 0x00, 0x00, 0x01],前三个字节都是UTF-8编码中的空字符(0x00),解码时会被截断,同样无可见内容。 - 第三种转字符串的写法:数值1转字符串后得到字符'1',对应UTF-8编码字节为
0x31,属于可打印ASCII字符,因此kafkacat可以正常显示。
正确实现方案
根据业务场景二选一即可,不存在所谓“不规范”的问题:
场景1:Header仅由业务代码按数值类型消费
如果消费端是Scala/Java业务代码,会直接将字节数组反序列化为Int/Long类型,不需要通过命令行工具直接查看明文值,ByteBuffer的写法就是标准实现,完全符合Kafka生态的二进制存储规范:
// 写入Int,ByteBuffer默认大端序和Kafka默认字节序一致,无需额外调整 val intBytes = ByteBuffer.allocate(4).putInt(counter).array() context.headers().add("counter_int", intBytes) // 消费端读取Int val readInt = ByteBuffer.wrap(record.headers().lastHeader("counter_int").value()).getInt() // 写入Long val longVal: Long = 100L val longBytes = ByteBuffer.allocate(8).putLong(longVal).array() context.headers().add("counter_long", longBytes) // 消费端读取Long val readLong = ByteBuffer.wrap(record.headers().lastHeader("counter_long").value()).getLong()
场景2:Header需要支持命令行工具直接查看明文
如果日常排查问题需要用kafkacat等运维工具直接看到Header的明文值,转UTF-8字符串存储的写法是完全合理的——Kafka生态中大量开源组件(Kafka Connect、Confluent Schema Registry相关组件)都采用字符串编码存储Header中的数值,只要上下游约定好编码格式就没有问题。可以直接使用Kafka自带的工具类简化转换逻辑:
import org.apache.kafka.common.utils.Utils // 写入 context.headers().add("counter_str", Utils.utf8(String.valueOf(counter))) // 读取 val readVal = Utils.utf8(record.headers().lastHeader("counter_str").value()).toInt
内容的提问来源于stack exchange,提问作者Jaime Caffarel
相关产品推荐
相关产品推荐

