Kafka AVRO序列化Joda DateTime类型报错的解决方法咨询
解决AVRO序列化Joda DateTime时的Unknown datum type错误
这个问题的核心原因很明确:你的AVRO Schema定义的是timestamp-millis逻辑类型(底层存储为long),但生成的Java类用了Joda DateTime,而AVRO默认的序列化器只认识原生long类型,不知道如何处理Joda DateTime对象,所以抛出了这个错误。下面是几种可行的解决办法:
方法一:手动转换为毫秒数(最简便)
既然timestamp-millis本质是存储时间的毫秒数,你只需要在设置字段时,把Joda DateTime转换成对应的long值即可,AVRO就能正常序列化:
// 假设你的生成类是ErrorRecord DateTime errorDateTime = new DateTime("2018-05-02T12:32:27.563Z"); ErrorRecord record = ErrorRecord.newBuilder() .setErrortime(errorDateTime.getMillis()) // 将DateTime转为毫秒数 .build(); // 执行序列化操作 DatumWriter<ErrorRecord> writer = new SpecificDatumWriter<>(ErrorRecord.class); ByteArrayOutputStream out = new ByteArrayOutputStream(); Encoder encoder = EncoderFactory.get().binaryEncoder(out, null); writer.write(record, encoder); encoder.flush();
方法二:自定义AVRO逻辑类型处理器(适合长期复用)
如果不想每次都手动转换,可以扩展AVRO的逻辑类型系统,让它直接支持Joda DateTime的序列化和反序列化:
- 创建一个自定义的
LogicalType实现,负责DateTime和long的转换:
public class JodaTimestampMillisLogicalType extends LogicalType { public static final String NAME = "timestamp-millis"; public JodaTimestampMillisLogicalType() { super(NAME); } @Override public void validate(Schema schema) { super.validate(schema); if (schema.getType() != Schema.Type.LONG) { throw new IllegalArgumentException("Joda timestamp-millis logical type must be backed by long"); } } @Override public Object deserialize(Object datum) { if (datum == null) return null; return new DateTime((Long) datum, DateTimeZone.UTC); } @Override public Object serialize(Object datum) { if (datum == null) return null; if (!(datum instanceof DateTime)) { throw new AvroRuntimeException("Unknown datum type " + datum.getClass()); } return ((DateTime) datum).getMillis(); } }
- 注册这个逻辑类型到AVRO的工厂:
LogicalTypes.register(JodaTimestampMillisLogicalType.NAME, new LogicalTypeFactory() { @Override public LogicalType fromSchema(Schema schema) { return new JodaTimestampMillisLogicalType(); } });
之后AVRO就能自动处理Joda DateTime的序列化和反序列化了。
方法三:改用Java 8时间类型(推荐,更现代)
如果你的项目已经支持Java 8+,可以调整AVRO代码生成的配置,让timestamp-millis对应Java 8的Instant类型(AVRO原生支持),这样就不需要额外处理:
以Maven插件为例,修改avro-maven-plugin的配置:
<plugin> <groupId>org.apache.avro</groupId> <artifactId>avro-maven-plugin</artifactId> <version>1.11.0</version> <!-- 使用较新版本的AVRO --> <executions> <execution> <phase>generate-sources</phase> <goals> <goal>schema</goal> </goals> <configuration> <logicalTypeImplementations> <logicalTypeImplementation> <logicalType>timestamp-millis</logicalType> <implementation>java.time.Instant</implementation> </logicalTypeImplementation> </logicalTypeImplementations> </configuration> </execution> </executions> </plugin>
重新生成Java类后,errortime字段会变成Instant类型,AVRO可以直接序列化,不需要额外转换。
内容的提问来源于stack exchange,提问作者Mikhail
相关产品推荐
相关产品推荐

