【问题标题】:Sending data from kafka producer in spark using scala使用scala在火花中从kafka生产者发送数据
【发布时间】:2014-10-12 02:03:29
【问题描述】:

在这段代码中:

import java.limport java.lang.StringBuilder
import java.util.Properties
import kafka.producer.{KeyedMessage, Producer, ProducerConfig}
import org.jnetpcap.Pcap
import org.jnetpcap.packet.{PcapPacket, PcapPacketHandler}

object kafkaproducer extends Serializable{
  def main(args: Array[String]) {
    if (args.length < 4) {
      System.err.println("Usage: KafkaWordCountProducer <metadataBrokerList> <topic> " +
        "<messagesPerSec> <wordsPerMessage>")
      System.exit(1)
    }
    //metadata.broker.list=localhost:9092
    //zookeeper.connect=localhost:2181
    val Array(brokers, topic, messagesPerSec, wordsPerMessage) = args
    // Zookeeper connection properties
    val props = new Properties()
    props.put("metadata.broker.list", brokers.toString)
    props.put("serializer.class", "kafka.serializer.StringEncoder")
    val config = new ProducerConfig(props)
    val producer = new Producer[String, PcapPacket](config)
    // Send some messages
    val snaplen = 64 * 1024 // Capture all packets, no truncation
    val flags = Pcap.MODE_PROMISCUOUS // capture all packets
    val timeout = 10 * 1000
    val jsb = new java.lang.StringBuilder()
    val errbuf = new StringBuilder(jsb);
    val pcap = Pcap.openLive("eth0", snaplen, flags, timeout, errbuf)
    if (pcap == null) {
      println("Error : " + errbuf.toString())
    }

    while(true){

      val jpacketHandler = new PcapPacketHandler[String]() {

        def nextPacket(packet: PcapPacket, user: String) {
          val data = new KeyedMessage[String,PcapPacket](topic.toString,(packet))
          println(data)
          producer.send(data)


        }
      }
      pcap.loop(50, jpacketHandler, "jNetPcap works!")


    }

  }
}

我得到了这个例外:

Exception in thread "main" java.lang.ClassCastException: org.jnetpcap.packet.PcapPacket cannot be cast to java.lang.String
at kafka.serializer.StringEncoder.toBytes(Unknown Source)
at kafka.producer.async.DefaultEventHandler$$anonfun$serialize$1.apply(Unknown Source)
at kafka.producer.async.DefaultEventHandler$$anonfun$serialize$1.apply(Unknown Source)
at scala.collection.TraversableLike$$anonfun$map$1.apply(TraversableLike.scala:244)
at scala.collection.TraversableLike$$anonfun$map$1.apply(TraversableLike.scala:244)
at scala.collection.IndexedSeqOptimized$class.foreach(IndexedSeqOptimized.scala:33)
at scala.collection.mutable.WrappedArray.foreach(WrappedArray.scala:34)
at scala.collection.TraversableLike$class.map(TraversableLike.scala:244)
at scala.collection.AbstractTraversable.map(Traversable.scala:105)
at kafka.producer.async.DefaultEventHandler.serialize(Unknown Source)
at kafka.producer.async.DefaultEventHandler.handle(Unknown Source)
at kafka.producer.Producer.send(Unknown Source)
at kafkaproducer$$anon$1.nextPacket(kafkaproducer.scala:50)
at kafkaproducer$$anon$1.nextPacket(kafkaproducer.scala:40)
at org.jnetpcap.Pcap.loop(Native Method)
at org.jnetpcap.Pcap.loop(Unknown Source)
at kafkaproducer$.main(kafkaproducer.scala:55)
at kafkaproducer.main(kafkaproducer.scala)
at sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method)
at sun.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:57)
at sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43)
at java.lang.reflect.Method.invoke(Method.java:606)
at com.intellij.rt.execution.application.AppMain.main(AppMain.java:134)

从 kafka 生产者发送数据时发生以下错误。在 kafka 生产者中,使用 jnetpcap 库捕获数据包。谁能帮帮我?

【问题讨论】:

标签: scala apache-spark apache-kafka


【解决方案1】:

这个异常的原因是这里生产者被配置为使用StringEncoder:

props.put("serializer.class", "kafka.serializer.StringEncoder")

尽管如此,提供的实际值是PcapPacket 类型。生产者将使用编码器序列化对象,boem 你有那个类转换异常。

还请注意,按照 JNetPcap 的文档,您可能不会使用捕获的 PcapPacket 来传输数据。该对象是可变的,并且在每次捕获时都会随着新捕获的数据而变化。 From the docs:

就像使用 JBufferHandler 一样,重复使用 PcapPacket 的单个副本 对于来自同一 pcap 调度循环实例的每个数据包。这 数据包到达时已完全解码,可以立即访问,但可以 不得放入队列或其他永久/半永久 存储。 需要立即由用户的 应用程序,丢弃或复制到更永久的内存位置。

正如我在this question 上提到的:

如果您想访问 PcapPacket 的详细信息,我建议 yIf 您 想要访问 PcapPacket 的细节,我建议你提取它 生产者端的信息并将其放入字符串或自定义中 可序列化的对象。

对于这种情况,这仍然是有效的建议。

【讨论】:

  • 谢谢。在页面下方的上述链接中,他们说可以将数据包放入队列中。那么我可以将相同的内容放在 kafka 队列中吗?
  • 我已经用如何实现编码器/解码器功能的指针更新了这个问题 - stackoverflow.com/questions/26258553/…
猜你喜欢
  • 2018-04-19
  • 2021-04-23
  • 2016-12-16
  • 2019-02-02
  • 1970-01-01
  • 2021-07-16
  • 2019-08-10
  • 2016-01-01
  • 1970-01-01
相关资源
最近更新 更多