如何在Kafka StateStore中存储含头部与时间戳的Record/ConsumerRecord?
问题解答
先搞懂 Record 和 ConsumerRecord 的区别
ConsumerRecord是Kafka Consumer API的原生类,包含消费到的全量元数据(topic、partition、offset、headers、timestamp等),但不适合存入StateStore——它的序列化逻辑依赖Kafka内部实现,且包含offset、partition这类和存储无关的字段,会额外增加存储开销。Record<K, V>是Kafka Streams提供的轻量级封装(属于org.apache.kafka.streams.processor包),仅保留key、value、timestamp、headers这几个你需要的核心字段,是专为Streams场景设计的,更适合作为StateStore的存储类型。
Kafka是否原生支持存储Record的Serde?
Kafka Streams没有提供原生Serde实现来直接序列化/反序列化Record,因为Record的泛型参数(K、V)是自定义的,框架无法提前适配所有类型的key和value。所以自定义Serde是目前唯一可行的方案,但实现逻辑并不复杂。
自定义Serde的Kotlin实现示例
下面是针对Record<String, CustomAVROSchema>的Serde实现,核心是复用你已有的StringSerde和SpecificAvroSerde处理key和value,同时单独序列化headers和timestamp:
1. 自定义Serializer
import org.apache.kafka.common.header.Headers import org.apache.kafka.common.header.internals.RecordHeaders import org.apache.kafka.common.serialization.Serializer import org.apache.kafka.streams.processor.Record import java.io.ByteArrayOutputStream import java.io.DataOutputStream class RecordAvroSerializer( private val keySerializer: Serializer<String>, private val valueSerializer: Serializer<CustomAVROSchema> ) : Serializer<Record<String, CustomAVROSchema>> { override fun serialize(topic: String?, data: Record<String, CustomAVROSchema>?): ByteArray? { if (data == null) return null ByteArrayOutputStream().use { baos -> DataOutputStream(baos).use { dos -> // 序列化timestamp(8字节长整型) dos.writeLong(data.timestamp()) // 序列化headers数量 dos.writeInt(data.headers().size()) // 逐个序列化header的key和value data.headers().forEach { header -> dos.writeUTF(header.key()) val value = header.value() if (value == null) { dos.writeInt(-1) } else { dos.writeInt(value.size) dos.write(value) } } // 序列化key val keyBytes = keySerializer.serialize(topic, data.key()) if (keyBytes == null) { dos.writeInt(-1) } else { dos.writeInt(keyBytes.size) dos.write(keyBytes) } // 序列化value(你的CustomAVROSchema) val valueBytes = valueSerializer.serialize(topic, data.value()) if (valueBytes == null) { dos.writeInt(-1) } else { dos.writeInt(valueBytes.size) dos.write(valueBytes) } return baos.toByteArray() } } } override fun configure(configs: MutableMap<String, *>?, isKey: Boolean) { keySerializer.configure(configs, isKey) valueSerializer.configure(configs, false) } override fun close() { keySerializer.close() valueSerializer.close() } }
2. 自定义Deserializer
import org.apache.kafka.common.header.Header import org.apache.kafka.common.header.internals.RecordHeader import org.apache.kafka.common.header.internals.RecordHeaders import org.apache.kafka.common.serialization.Deserializer import org.apache.kafka.streams.processor.Record import java.io.ByteArrayInputStream import java.io.DataInputStream class RecordAvroDeserializer( private val keyDeserializer: Deserializer<String>, private val valueDeserializer: Deserializer<CustomAVROSchema> ) : Deserializer<Record<String, CustomAVROSchema>> { override fun deserialize(topic: String?, data: ByteArray?): Record<String, CustomAVROSchema>? { if (data == null) return null ByteArrayInputStream(data).use { bais -> DataInputStream(bais).use { dis -> // 反序列化timestamp val timestamp = dis.readLong() // 反序列化headers val headerCount = dis.readInt() val headers = RecordHeaders() repeat(headerCount) { val key = dis.readUTF() val valueLength = dis.readInt() val value = if (valueLength == -1) null else { val bytes = ByteArray(valueLength) dis.readFully(bytes) bytes } headers.add(RecordHeader(key, value)) } // 反序列化key val keyLength = dis.readInt() val key = if (keyLength == -1) null else { val keyBytes = ByteArray(keyLength) dis.readFully(keyBytes) keyDeserializer.deserialize(topic, keyBytes) } // 反序列化value val valueLength = dis.readInt() val value = if (valueLength == -1) null else { val valueBytes = ByteArray(valueLength) dis.readFully(valueBytes) valueDeserializer.deserialize(topic, valueBytes) } return Record(key, value, timestamp, headers) } } } override fun configure(configs: MutableMap<String, *>?, isKey: Boolean) { keyDeserializer.configure(configs, isKey) valueDeserializer.configure(configs, false) } override fun close() { keyDeserializer.close() valueDeserializer.close() } }
3. 封装成Serde
import org.apache.kafka.common.serialization.Serde import org.apache.kafka.common.serialization.Serdes import org.apache.kafka.streams.processor.Record class RecordAvroSerde : Serde<Record<String, CustomAVROSchema>> { private val innerSerde: Serde<Record<String, CustomAVROSchema>> init { val keySerde = Serdes.String() val valueSerde = SpecificAvroSerde<CustomAVROSchema>().apply { // 传入你原有的Avro配置,比如schema registry地址 configure(mapOf( "schema.registry.url" to "http://your-schema-registry:8081" ), false) } innerSerde = Serdes.serdeFrom( RecordAvroSerializer(keySerde.serializer(), valueSerde.serializer()), RecordAvroDeserializer(keySerde.deserializer(), valueSerde.deserializer()) ) } override fun serializer() = innerSerde.serializer() override fun deserializer() = innerSerde.deserializer() override fun configure(configs: MutableMap<String, *>?, isKey: Boolean) { innerSerde.configure(configs, isKey) } override fun close() { innerSerde.close() } }
使用自定义Serde创建StateStore
在构建PersistentTimestampedKeyValueStore时,指定你的自定义Serde即可:
val storeBuilder = Stores.keyValueStoreBuilder( Stores.persistentTimestampedKeyValueStore("your-store-name"), Serdes.String(), RecordAvroSerde() ) streamsBuilder.addStateStore(storeBuilder)
额外注意事项
- 如果你不需要
Record的全量字段,可以自定义一个更精简的数据类(仅包含key、value、timestamp、headers),能进一步降低存储开销。 - 确保
CustomAVROSchema的Serde配置和原有逻辑一致,避免序列化/反序列化失败。 - 测试时可通过Streams的
InteractiveQuery验证存储内容是否符合预期。
内容的提问来源于stack exchange,提问作者HaavardG
相关产品推荐
相关产品推荐

