【问题标题】:How to implement Generic Kafka Streams Deserializer如何实现通用 Kafka Streams Deserializer
【发布时间】:2018-05-24 13:58:26
【问题描述】:

我喜欢 Kafka,但讨厌编写大量的序列化器/反序列化器,所以我尝试创建一个可以反序列化泛型 T 的GenericDeserializer<T>

这是我的尝试:

class GenericDeserializer< T > implements Deserializer< T > {
    static final ObjectMapper objectMapper = new ObjectMapper();
    @Override
    public void configure(Map<String, ?> configs, boolean isKey) {
    }
    @Override
    public T deserialize( String topic, byte[] data) {
            T result = null;
            try {
                    result = ( T )( objectMapper.readValue( data, T.class ) );
            }
            catch ( Exception e ) {
                    e.printStackTrace();
            }
            return result;
    }
    @Override
    public void close() {
    }
}

但是,(Eclipse)Java 编译器抱怨行

result = ( T )( objectMapper.readValue( data, T.class ) );

带有消息Illegal class literal for the type parameter T

问题:

  1. 您能解释一下消息的含义吗?
  2. 有没有办法解决这个问题以获得预期的效果?

【问题讨论】:

  • 问题是,泛型类型只能在编译时用于类型检查。在运行时,所有Ts 都替换为Object 类型。因此,T.class 无法评估...它被称为“类型擦除”,在 Java 中非常烦人...

标签: generic-programming apache-kafka-streams


【解决方案1】:

您可以使用包com.fasterxml.jackson.core.type中的TypeReference 实现通用反序列化

public class KafkaGenericDeserializer<T> implements Deserializer<T> {

    private final ObjectMapper mapper;
    private final TypeReference<T> typeReference;

    public KafkaGenericDeserializer(ObjectMapper mapper, TypeReference<T> typeReference) {
        this.mapper = mapper;
        this.typeReference = typeReference;
    }

    @Override
    public T deserialize(final String topic, final byte[] data) {
        if (data == null) {
            return null;
        }

        try {
            return mapper.readValue(data, typeReference);
        } catch (final IOException ex) {
            throw new SerializationException("Can't deserialize data [" + Arrays.toString(data) + "] from topic [" + topic + "]", ex);
        }
    }

    @Override
    public void close() {}

    @Override
    public void configure(final Map<String, ?> settings, final boolean isKey) {}
}

使用这样的通用反序列化器,你可以创建Serge:

public static <T> Serde<T> createSerdeWithGenericDeserializer(TypeReference<T> typeReference) {
    KafkaGenericDeserializer<T> kafkaGenericDeserializer = new KafkaGenericDeserializer<>(objectMapper, typeReference);
    return Serdes.serdeFrom(new JsonSerializer<>(), kafkaGenericDeserializer);
}

这里的JsonSerializer是来自spring-kafka的依赖,或者实现自己的序列化。

之后,您可以在 Kafka 流创建期间使用 serde:

TypeReference<YourGenericClass<SpecificClass>> typeReference = new TypeReference<YourGenericClass<SpecificClass>>() {};
Serde<YourGenericClass<SpecificClass>> itemSerde = createSerdeWithGenericDeserializer(typeReference);
Consumed<String, YourGenericClass<SpecificClass>> consumed = Consumed.with(Serdes.String(), itemSerde);
streamsBuilder.stream(topicName, consumed);

【讨论】:

  • 这里你传入的对象映射器是什么 createSerdeWithGenericDeserializer 方法 "new KafkaGenericDeserializer(objectMapper, typeReference);"
  • 其实这取决于你对反序列化的要求。您可以使用 Spring 的实现 org.springframework.http.converter.json.Jackson2ObjectMapperFactoryBean: ObjectMapper objectMapper = Jackson2ObjectMapperBuilder.json().featuresToDisable(FAIL_ON_UNKNOWN_PROPERTIES).serializationInclusion(NON_NULL).build(); 或只是 new ObjectMapper() 并根据您的需要进行配置
【解决方案2】:

在java中,你不能实例化一个泛型类型,即使是反射性的,这意味着objectMapper.readValue()不能用T.class做任何事情。因此,您需要知道在给定情况下要创建什么类。这样做的合乎逻辑的方法是有一些主题 -> 类型的映射,您的反序列化程序可以访问。这方面的一个例子是SpecificAvroSerde,它使用融合模式注册表(一个外部进程)来识别要反序列化的类型。您也可以将此映射构建到您的代码中,但根据您的用例,这不会特别健壮。

https://github.com/confluentinc/schema-registry/blob/master/avro-serde/src/main/java/io/confluent/kafka/streams/serdes/avro/SpecificAvroSerde.java

SpecificAvroSerde 的内容更深一些 - 这里有一个块正在执行询问模式注册表应该解码为什么类型的工作: https://github.com/confluentinc/schema-registry/blob/master/avro-serializer/src/main/java/io/confluent/kafka/serializers/AbstractKafkaAvroDeserializer.java#L109-L139

当然,这段代码都被 Avro 的复杂性所笼罩。如果有时间,我会编写一些示例代码,说明如何使用 JSON 在内存中执行此操作。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2022-08-24
    • 2018-07-29
    • 1970-01-01
    • 2019-10-11
    • 1970-01-01
    • 2020-06-02
    • 1970-01-01
    • 2020-07-30
    相关资源
    最近更新 更多