【问题标题】:How to consume JSON records from Kafka using Spark Streaming and Python?如何使用 Spark Streaming 和 Python 使用来自 Kafka 的 JSON 记录?
【发布时间】:2017-10-24 18:22:11
【问题描述】:

我创建了一个带有 JSON 格式记录的 Kafka 主题。

我可以使用 kafka-console-consumer.sh 来使用这些 JSON 字符串:

./kafka-console-consumer.sh --new-consumer \
    --topic test \
    --from-beginning \
    --bootstrap-server host:9092 \
    --consumer.config /root/client.properties

如何在 Python 中使用 Spark Streaming 来做到这一点?

【问题讨论】:

    标签: python apache-spark apache-kafka spark-streaming


    【解决方案1】:

    Doh,为什么 Python 不是 Scala?! 那么你的家庭练习就是将下面的代码重写为 Python ;-)

    来自Advanced Sources:

    从 Spark 2.1.1 开始,在这些源中,Kafka、Kinesis 和 Flume 在 Python API 中可用。

    基本上,这个过程是:

    使用spark-streaming-kafka-0-10_2.11 库从Kafka 主题读取消息,如Spark Streaming + Kafka Integration Guide (Kafka broker version 0.10.0 or higher) 中所述,使用KafkaUtils.createDirectStream。

    import org.apache.kafka.clients.consumer.ConsumerRecord
    import org.apache.kafka.common.serialization.StringDeserializer
    import org.apache.spark.streaming.kafka010._
    import org.apache.spark.streaming.kafka010.LocationStrategies.PreferConsistent
    import org.apache.spark.streaming.kafka010.ConsumerStrategies.Subscribe
    
    val kafkaParams = Map[String, Object](
      "bootstrap.servers" -> "localhost:9092,anotherhost:9092",
      "key.deserializer" -> classOf[StringDeserializer],
      "value.deserializer" -> classOf[StringDeserializer],
      "group.id" -> "use_a_separate_group_id_for_each_stream",
      "auto.offset.reset" -> "latest",
      "enable.auto.commit" -> (false: java.lang.Boolean)
    )
    
    val topics = Array("topicA", "topicB")
    val stream = KafkaUtils.createDirectStream[String, String](
      streamingContext,
      PreferConsistent,
      Subscribe[String, String](topics, kafkaParams)
    )
    

    使用 map 运算符将 ConsumerRecords 复制到值,这样您就不会遇到序列化问题。

    stream.map(record => (record.key, record.value))
    

    如果不发送密钥,record.value 就足够了。

    stream.map(record => record.value)
    

    将字符串消息转换为 JSON 获得值后,使用 from_json 函数:

    from_json(e: Column, schema: StructType) 将包含 JSON 字符串的列解析为具有指定架构的 StructType。在不可解析字符串的情况下返回null。

    代码如下:

    ...foreach { rdd =>
      messagesRDD.toDF.
        withColumn("json", from_json('value, jsonSchema)).
        select("json.*").show(false)
    }
    

    完成!

    【讨论】:

    • 非常感谢您抽出宝贵时间提供答案,我会试一试并告诉您。
    猜你喜欢
    • 1970-01-01
    • 2016-11-03
    • 1970-01-01
    • 2020-03-08
    • 2018-01-09
    • 2021-04-03
    • 2021-11-19
    • 1970-01-01
    • 2019-04-03
    相关资源
    最近更新 更多