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)
}
完成!