【问题标题】:Kafka - consuming from spark卡夫卡——从火花中消费
【发布时间】:2018-04-25 22:36:54
【问题描述】:

我关注了this document,效果很好。现在我尝试使用 spark 中的连接器数据。有什么可以参考的吗?由于我使用的是confluent,它与原始kafka参考文档有很大不同。

这是我目前使用的一些代码。问题是无法将记录数据转换为 java.String。 (并且不确定这是正确的消费方式)

val brokers = "http://127.0.0.1:9092"
val topics = List("postgres-accounts2")
val sparkConf = new SparkConf().setAppName("KafkaWordCount")
//sparkConf.setMaster("spark://sda1:7077,sda2:7077")
sparkConf.setMaster("local[2]")
sparkConf.set("spark.serializer", "org.apache.spark.serializer.KryoSerializer") 
sparkConf.registerKryoClasses(Array(classOf[org.apache.avro.generic.GenericData$Record]))

val ssc = new StreamingContext(sparkConf, Seconds(2))
ssc.checkpoint("checkpoint")


 // Create direct kafka stream with brokers and topics
//val topicsSet = topics.split(",")

val kafkaParams = Map[String, Object](
  "schema.registry.url" -> "http://127.0.0.1:8081",
  "bootstrap.servers" -> "http://127.0.0.1:9092",
  "key.deserializer" -> "io.confluent.kafka.serializers.KafkaAvroDeserializer",
   "value.deserializer" -> "io.confluent.kafka.serializers.KafkaAvroDeserializer",
  "group.id" -> "use_a_separate_group_id_for_each_stream",
  "auto.offset.reset" -> "earliest",
  "enable.auto.commit" -> (false: java.lang.Boolean)
)

val messages = KafkaUtils.createDirectStream[String, String](
  ssc,
  PreferConsistent,
  Subscribe[String, String](topics, kafkaParams)
)

val data = messages.map(record => {
    println( record) 
    println( "value : " + record.value().toString() ) // error  java.lang.ClassCastException: org.apache.avro.generic.GenericData$Record cannot be cast to java.lang.String
    //println( Json.parse( record.value() + ""))

    (record.key, record.value)
})

【问题讨论】:

  • 您似乎刚刚开始使用 Spark 及其流功能,让我问您为什么使用 Spark Streaming 而不是 Spark Structured Streaming?见spark.apache.org/docs/latest/…
  • 由于我使用了confluent的jdbc连接器,我认为它应该是结构化的straming。不是吗?
  • 不知道 Confluent 中的 jdbc 连接器,但我相信您应该使用 Spark Structured Streaming 作为 Spark 中的流解决方案。
  • 以数据为记录类型解决了:)
  • 你能回答你自己的问题并批准吗?我期待看到变化。

标签: apache-spark apache-kafka avro confluent-schema-registry


【解决方案1】:

将我的值反序列化器同步到下面。它将提供适当的功能和类型。

KafkaUtils.createDirectStream[String, record]

【讨论】:

猜你喜欢
  • 2017-06-29
  • 2016-08-03
  • 2018-09-15
  • 2018-02-24
  • 2023-03-19
  • 2018-08-13
  • 2018-08-15
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多