【发布时间】:2016-12-14 23:31:37
【问题描述】:
我正在使用 Spark Streaming 从 Kafka 获取批量 JSON 读数。生成的批次从 RDD 转换为数据帧。
我的目标是对此数据帧的每一行进行分类,因此我使用 VectorAssembler 创建将传递给模型的特征:
sqlContext = SQLContext(rdd.context)
rawReading = sqlContext.jsonRDD(rdd)
sensorReadings = rawReading.selectExpr("actual.y AS yaw","actual.p AS pitch", "actual.r AS roll")
assembler = VectorAssembler(
inputCols=["yaw", "pitch", "roll"],
outputCol="features")
sensorReadingsFinal = assembler.transform(sensorReadings)
sensorReadingsFinal.show()
+---+-----+----+-----------------+
|yaw|pitch|roll| features|
+---+-----+----+-----------------+
| 18| 17.5| 120|[18.0,17.5,120.0]|
| 18| 17.5| 120|[18.0,17.5,120.0]|
| 18| 17.5| 120|[18.0,17.5,120.0]|
| 18| 17.5| 120|[18.0,17.5,120.0]|
| 18| 17.5| 120|[18.0,17.5,120.0]|
+---+-----+----+-----------------+
我有一个我之前训练过的随机森林模型。
loadedModel = RandomForestModel.load(sc, "MyRandomForest.model")
我的问题是,在将整个数据插入数据库之前,如何对数据框中的每一行进行预测?
我最初正在考虑做这样的事情......
prediction = loadedModel.predict(sensorReadings.features)
但我意识到,由于数据框有多行,我需要以某种方式添加一列并逐行进行预测。也许我在这一切都错了?
我想要的最终数据框是这样的:
+---+-----+----+-----------------+
|yaw|pitch|roll| Prediction|
+---+-----+----+-----------------+
| 18| 17.5| 120| 1 |
| 18| 17.5| 120| 1 |
| 18| 17.5| 120| 1 |
| 18| 17.5| 120| 1 |
| 18| 17.5| 120| 1 |
+---+-----+----+-----------------+
此时我会将其保存到数据库中:
sensorReadingsFinal.write.jdbc("jdbc:mysql://localhost/testdb", "SensorReadings", properties=connectionProperties)
【问题讨论】:
标签: python-2.7 apache-spark spark-streaming spark-dataframe apache-spark-mllib