【问题标题】:Using multiple deserializers for a kafka consumer为 kafka 消费者使用多个反序列化器
【发布时间】:2017-03-06 06:03:52
【问题描述】:

我是 kafka 甚至序列化的新手。到目前为止,我需要处理使用简单代码序列化的 json 格式的 kafka 事件。但现在正在使用 Avro 编码器添加额外的事件。所以现在我希望这个单一的消费者在 json 中使用 StringDeserialzer,对于 Avro 使用其各自的反序列化器。但是如何在同一个属性文件中映射 2 个反序列化器?

private Properties getProps(){
    Properties props = new Properties();
    props.put("group.id", env.getProperty("group.id"));
    props.put("enable.auto.commit", env.getProperty("enable.auto.commit"));
    props.put("key.deserializer", env.getProperty("key.deserializer"));
    props.put("value.deserializer", env.getProperty("value.deserializer"));
    return props;
}//here as only value can be mapped to "key.deserializer" is there anyway to do this

在主方法中

KafkaConsumer<String, String> _consumer = new KafkaConsumer<>(getProps());
consumers.add(_consumer);
_consumer.subscribe(new ArrayList<>(topicConsumer.keySet()));

【问题讨论】:

标签: java serialization properties avro kafka-consumer-api


【解决方案1】:

只需编写一个通用的反序列化器,它将主题委托给匹配的反序列化器。

public class GenericDeserializer extends JsonDeserializer<Object>
{
    public GenericDeserializer()
    {
    }

    @Override
    public Object deserialize(String topic, Headers headers, byte[] data)
    {
        switch (topic)
        {
        case KafkaTopics.TOPIC_ONE:
            TopicOneDeserializer topicOneDeserializer = new TopicOneDeserializer();
            topicOneDeserializer.addTrustedPackages("com.xyz");
            return topicOneDeserializer.deserialize(topic, headers, data);
        case KafkaTopics.TOPIC_TWO:
            TopicTwoDeserializer topicTwoDeserializer= new TopicTwoDeserializer();
            topicTwoDeserializer.addTrustedPackages("com.xyz");
            return topicTwoDeserializer.deserialize(topic, headers, data);
        }
        return super.deserialize(topic, data);
    }
}

【讨论】:

    【解决方案2】:

    您需要提供一个包含两个原始解串器的混合解串器。在内部,新的包装反序列化器必须能够区分这两种类型的消息,并将原始字节转发到执行实际工作的正确反序列化器。

    如果你不能提前知道你有什么类型的消息,你也可以尝试一个错误的方法——即默认情况下将它交给一个序列化程序,如果这个失败(即抛出异常)尝试第二个一个。

    【讨论】:

    • 我是反序列化的新手,Avro 反序列化被证明是困难的,因为这需要 pojo 并返回对象,但 StringDeserializer 返回字符串,我使用消费者中的 Gson 方法将其转换为对象。如何创建一个通用的方法来使用 try 和 catch,因为它们的返回类型不同?
    • 你还需要一个混合返回类型;基本上是一个包装器类,其中包含您包装的每种类型的成员。并且每次您创建混合类型的消息时,除了一个成员之外,所有成员都将是null。如何编写序列化程序显示在对您问题的评论中的链接中。
    猜你喜欢
    • 2019-08-30
    • 2019-04-29
    • 2018-12-08
    • 1970-01-01
    • 1970-01-01
    • 2019-05-11
    • 2021-10-06
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多