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

如何为KCache创建List<MyAvroObject>集合对应的Avro Serdes实例

解决方案

你要实现List<MyAvroObject>类型的Serde非常简单,直接用Kafka官方提供的集合Serde包装你已经写好的单个Avro对象Serde即可。

你之前的test变量写法无效是因为Serdes.serdeFrom仅支持Kafka内置的基础类型序列化/反序列化逻辑,无法识别自定义Avro类型的集合结构。

正确实现代码如下:

class Cache() {
    // 改成你期望的List类型声明
    private val cache: KafkaCache<String, List<MyAvroObject>>

    init {
        val cacheProps = Properties()
        cacheProps[KafkaCacheConfig.KAFKACACHE_BOOTSTRAP_SERVERS_CONFIG] =
            appConfig.propertyOrNull("kafka.bootstrapServers")?.getString()
        cacheProps[KafkaCacheConfig.KAFKACACHE_GROUP_ID_CONFIG] =
            appConfig.propertyOrNull("applicationId")?.getString()
        cacheProps[KafkaCacheConfig.KAFKACACHE_CLIENT_ID_CONFIG] = "GeneratingUnitsCache"

        val serdeConfig = Collections.singletonMap(
            AbstractKafkaSchemaSerDeConfig.SCHEMA_REGISTRY_URL_CONFIG,
            appConfig.property("kafka.schemaRegistryUrl").getString(),
        )

        // 先初始化单个MyAvroObject的Avro Serde
        val elementSerde: Serde<MyAvroObject> = SpecificAvroSerde()
        elementSerde.configure(serdeConfig, false)
        // 包装成List<MyAvroObject>的Serde,Kafka 2.6+ 自带ListSerde
        val listAvroSerde: Serde<List<MyAvroObject>> = ListSerde(elementSerde)

        cache = KafkaCache(
            KafkaCacheConfig(cacheProps),
            Serdes.String(),
            // 传入生成的List类型Serde
            listAvroSerde
        )

        cache.init()
    }
    
    // 其余类实现逻辑
}

如果你的Kafka版本低于2.6没有内置ListSerde,也可以手动简单实现一个包装类,分别实现序列化和反序列化逻辑遍历处理每个元素即可,不需要额外修改Avro schema相关配置,序列化后对应的Avro结构是array类型,元素为MyAvroObject的schema,完全兼容schema registry的校验逻辑。

内容的提问来源于stack exchange,提问作者hermanjakobsen

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.03 18:54:05