【发布时间】:2017-07-14 10:31:28
【问题描述】:
我正在尝试处理来自 Kinesis 的 Json 字符串。 Json 字符串可以有几种不同的形式。从 Kinesis,我创建了一个 DStream:
val kinesisStream = KinesisUtils.createStream(
ssc, appName, "Kinesis_Stream", "kinesis.ap-southeast-1.amazonaws.com",
"region", InitialPositionInStream.LATEST, kinesisCheckpointInterval, StorageLevel.MEMORY_AND_DISK_2)
val lines = kinesisStream.map(x => new String(x))
lines.foreachRDD((rdd, time) =>{
val sqlContext = SQLContextSingleton.getInstance(rdd.sparkContext)
import sqlContext.implicits.StringToColumn
if(rdd.count() > 0){
// Process jsons here
// Json strings here would have either one of the formats below
}
})
RDD 字符串将具有这些 json 字符串之一。 收藏:
[
{
"data": {
"ApplicationVersion": "1.0.3 (65)",
"ProjectId": 30024,
"TargetId": "4138",
"Timestamp": 0
},
"host": "host1"
},
{
"data": {
"ApplicationVersion": "1.0.3 (65)",
"ProjectId": 30025,
"TargetId": "4139",
"Timestamp": 0
},
"host": "host1"
}
]
有些 Json 字符串是像这样的单个对象:
{
"ApplicationVersion": "1.0.3 (65)",
"ProjectId": 30026,
"TargetId": "4140",
"Timestamp": 0
}
我希望能够从“data”键中提取对象,如果它是第一种Json字符串并与第二种Json结合形成RDD/DataFrame,我该如何实现?
最终我希望我的数据框是这样的:
+------------------+---------+--------+---------+
|ApplicationVersion|ProjectId|TargetId|Timestamp|
+------------------+---------+--------+---------+
| 1.0.3 (65)| 30024| 4138| 0|
| 1.0.3 (65)| 30025| 4139| 0|
| 1.0.3 (65)| 30026| 4140| 0|
+------------------+---------+--------+---------+
抱歉,Scala 和 Spark 的新手。我一直在查看现有示例,但遗憾的是没有找到解决方案。
非常感谢。
【问题讨论】:
标签: json scala apache-spark