【发布时间】:2018-08-01 13:32:57
【问题描述】:
情况
我目前正在使用 AVRO 和模式存储库编写消费者/生产者。
根据我收集到的信息,我对这些数据进行序列化的选项是使用 Confluent 的 avro 序列化程序,或者使用 Twitter 的 Bijection。
似乎双射看起来最直接。
所以我想以以下格式生成日期ProducerRecord[String,Array[Byte]],这归结为 [一些字符串 ID,序列化的 GenericRecord]
(注意:我要使用通用记录,因为此代码库必须处理从 Json/csv/... 解析的数千个模式)
问题:
我序列化和使用 AVRO 的全部原因是您不需要在数据本身中包含架构(就像使用 Json/XML/...一样)。
然而,当检查主题中的数据时,我看到整个方案与数据一起包含。我做错了什么,这是设计使然,还是应该改用融合序列化程序?
代码:
def jsonStringToAvro(jString: String, schema: Schema): GenericRecord = {
val converter = new JsonAvroConverter
val genericRecord = converter.convertToGenericDataRecord(jString.replaceAll("\\\\/","_").getBytes(), schema)
genericRecord
}
def serializeAsByteArray(avroRecord: GenericRecord): Array[Byte] = {
//val genericRecordInjection = GenericAvroCodecs.toBinary(avroRecord.getSchema)
val r: Array[Byte] = GenericAvroCodecs.toBinary(avroRecord.getSchema).apply(avroRecord)
r
}
//schema comes from a rest call to the schema repository
new ProducerRecord[String, Array[Byte]](topic, myStringKeyGoesHere, serializeAsByteArray(jsonStringToAvro(jsonObjectAsStringGoesHere, schema)))
producer.send(producerRecord, new Callback {...})
【问题讨论】:
-
双射库不与模式注册表交互,并且您不会像 Confluent 序列化程序那样将 ID 放在任何地方。因此,整个架构将成为消息的一部分
-
另外,
ProducerRecord[String, GenericRecord]有什么问题?并将 REST 调用放入序列化程序中? -
一些项目可能有很大的容量,所以我认为序列化会给我带来一些性能提升。那么,当我使用 Confluent 序列化程序时,这会从通用记录中剥离模式吗?
-
Kafka Serializer 接口旨在为您获取字节数组。在类的“主要方法”中编写该逻辑没有任何好处
标签: scala apache-kafka avro bijection