【发布时间】:2018-03-02 19:34:19
【问题描述】:
我可以用Apache Avro 发送和接收日期类型吗?我无法找到任何关于此的内容。只有我发现的东西说在模式中使用日期的 int 和logicalType。但这会在接收方产生另一个 int 。我仍然需要将其转换为日期。
我正在尝试从Apache Kafka 生产者发送日期并在 Kafka 消费者中接收。
如果没有其他方法,那么我是否必须始终将日期转换为 int,然后再返回给消费者。有这篇文章展示了如何做到这一点:
Get the number of days, weeks, and months, since Epoch in Java
序列化代码:-
@Override
public byte[] serialize(String topic, T data) {
try {
byte[] result = null;
if (data != null) {
logger.debug("data='{}'" + data);
ByteArrayOutputStream byteArrayOutputStream = new ByteArrayOutputStream();
BinaryEncoder binaryEncoder =
EncoderFactory.get().binaryEncoder(byteArrayOutputStream, null);
DatumWriter<GenericRecord> datumWriter = new GenericDatumWriter<>(data.getSchema());
datumWriter.write(data, binaryEncoder);
binaryEncoder.flush();
byteArrayOutputStream.close();
result = byteArrayOutputStream.toByteArray();
byteArrayOutputStream.close();
logger.debug("serialized data='{}'" + DatatypeConverter.printHexBinary(result));
}
return result;
} catch (IOException ex) {
throw new SerializationException(
"Can't serialize data='" + data + "' for topic='" + topic + "'", ex);
}
}
解串器代码:-
@Override
public T deserialize(String topic, byte[] data) {
try {
T result = null;
if (data != null) {
logger.debug("data='{}'" + DatatypeConverter.printHexBinary(data));
DatumReader<GenericRecord> datumReader =
new SpecificDatumReader<>(targetType.newInstance().getSchema());
Decoder decoder = DecoderFactory.get().binaryDecoder(data, null);
result = (T) datumReader.read(null, decoder);
logger.debug("deserialized data='{}'" + result);
}
return result;
} catch (Exception ex) {
throw new SerializationException(
"Can't deserialize data '" + Arrays.toString(data) + "' from topic '" + topic + "'", ex);
}
}
架构文件:-
{"namespace": "com.test",
"type": "record",
"name": "Measures",
"fields": [
{"name": "transactionDate", "type": ["int", "null"], "logicalType" : "date" }
]
}
这两个只是在生产者和消费者配置中定义为序列化器和反序列化器类。
【问题讨论】:
-
“接收方的
int”是什么意思?您反序列化为的 Java 类型应该有一个 Avro 可以填充的Date字段。我也强烈建议不要使用Date- 如果您需要一个时间点,请使用Instant。 -
您的架构错误 - 逻辑类型在类型上而不是在字段上。
{ "type": "long", "logicalType": "date" } -
您有什么理由使用自己的解码器? Confluent 提供了自己的docs.confluent.io/current/schema-registry/docs/…
-
谢谢@BoristheSpider,我的架构是错误的,在更正它并在阅读下面的 Basil 回复后使用适配器后,我可以让它像 joda LocalDate 一样工作。我想尽可能避免融合,但没有任何理由。
标签: java apache-kafka avro