如何为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
相关产品推荐
相关产品推荐

