【问题标题】:How to implement deserialisation in kafka consumer using scala?如何使用scala在kafka消费者中实现反序列化?
【发布时间】:2014-11-02 17:47:47
【问题描述】:

我的 kafka 消费者代码中有以下行。

val lines = KafkaUtils.createStream(ssc, zkQuorum, group, topicpMap).map(_._2) 

如何将此流“行”反序列化为原始对象?通过将类扩展为可序列化,在 kafka 生产者中实现了可序列化。我正在使用 scala 在 spark 中实现这一点。

【问题讨论】:

    标签: scala deserialization apache-spark apache-kafka


    【解决方案1】:

    您需要实现自定义Decoder,并将预期的类型信息与解码器一起提供给createStream 函数。

    KafkaUtils.createStream[KeyType, ValueType, KeyDecoder, ValueDecoder] (...)

    例如,如果您使用String 作为键和CustomContainer 作为值,您的流创建将如下所示:

    val stream = KafkaUtils.createStream[String, CustomContainer, StringDecoder, CustomContainerDecoder](...)  
    

    鉴于您将消息编码为new KeyedMessage[String,String],因此正确的解码器是这样的字符串解码器:

    KafkaUtils.createStream[String, String, StringDecoder, StringDecoder](topic,...)
    

    这将为您提供DStream[String,String] 作为您处理的基础。

    如果你想发送/接收特定的对象类型,你需要为它实现一个 Kafka EncoderDecoder。 幸运的是,PcapPacket 已经实现了您需要的方法:

    剩下的就是实现 Kafka 所需的 Encoder/Decoder 接口的样板代码。

    【讨论】:

    • github.com/swe0523/producer/blob/master/producer.scala 这是我的生产者代码的链接。你能帮帮我吗?
    • @user3823859 查看更新的答案以及对您的具体问题的反馈。
    • 谢谢。在此之后,是否可以在 kafka 消费者中使用类似 stream.getCaptureHeader() 的功能,这些功能在 kafka 生产者中使用的 jnetpcap 库中可用?
    • 不,您将消息序列化为字符串。您必须解析字符串来处理数据。
    • 有没有内置函数可以做同样的事情?
    猜你喜欢
    • 2018-12-08
    • 2017-03-06
    • 2019-05-11
    • 2019-08-30
    • 1970-01-01
    • 2020-08-20
    • 2019-04-29
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多