您需要实现自定义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 Encoder 和 Decoder。
幸运的是,PcapPacket 已经实现了您需要的方法:
剩下的就是实现 Kafka 所需的 Encoder/Decoder 接口的样板代码。