【发布时间】:2021-07-22 21:41:09
【问题描述】:
我有 2 个架构:
Event.avsc:
{
"type": "record",
"namespace": "com.onemount.jobs.transform.schema.avro",
"name": "Event",
"fields": [
{
"name": "id",
"type": "string"
},
{
"name": "mtp_interest_submit",
"type": ["null", "InterestSubmitParam"],
"default": null
}
]
}
InterestSubmitParam.avsc:
{
"type": "record",
"namespace": "com.onemount.jobs.transform.schema.avro",
"name": "InterestSubmitParam",
"fields": [
{
"name": "interest",
"type": {
"type": "array",
"items": "string"
}
}
]
}
我正在使用来自 Kafka Confluent(使用 specific.avro.reader=false)的 Avro 消息,并且需要从 GenericRecord 转换为 ObjectNode。结果如下:
{
"id": "c8b76e58-9803-4c78-9f82-a185bda1cabf",
"mtp_interest_submit": {
"com.onemount.jobs.transform.schema.avro.InterestSubmitParam": {
"interest": [
"fashion",
"travel"
]
}
}
}
但我预计应该是:
{
"id": "c8b76e58-9803-4c78-9f82-a185bda1cabf",
"mtp_interest_submit": {
"interest": [
"fashion",
"travel"
]
}
}
我该如何解决。这是我的转换器代码:
GenericRecord genericRecord = ...
try (ByteArrayOutputStream outputStream = new ByteArrayOutputStream()) {
DatumWriter<GenericRecord> writer = new GenericDatumWriter<>(genericRecord.getSchema());
JsonEncoder encoder = EncoderFactory.get().jsonEncoder(genericRecord.getSchema(), outputStream);
writer.write(genericRecord, encoder);
encoder.flush();
return new String(outputStream.toByteArray(), StandardCharsets.UTF_8);
}
非常感谢!
【问题讨论】:
-
不要让
mtp_interest_submit可以为空并且类型名称不会在那里 -
@OneCricketeer ty。添加非空字段可能会违反完全兼容性,我将在最后考虑这种方法
-
另一种方法是构建一个新的 POJO 模型,但它实际上与从模式构建 SpecificRecord 的类定义相同,所以也许可以考虑使用 Jackson Avro Objectmapper 而不是普通的 Avro API跨度>
标签: java avro confluent-schema-registry fasterxml