You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.27 02:57:06