如何查询具有复杂类型(如地图/数组)的RDD?例如,当我写这个测试代码时:
case class Test(name: String, map: Map[String, String])
val map = Map("hello" -> "world", "hey" -> "there")
val map2 = Map("hello" -> "people", "hey" -> "you")
val rdd = sc.parallelize(Array(Test("first", map), Test("second", map2)))
Run Code Online (Sandbox Code Playgroud)
我虽然语法如下:
sqlContext.sql("SELECT * FROM rdd WHERE map.hello = world")
Run Code Online (Sandbox Code Playgroud)
要么
sqlContext.sql("SELECT * FROM rdd WHERE map[hello] = world")
Run Code Online (Sandbox Code Playgroud)
但我明白了
无法访问MapType类型中的嵌套字段(StringType,StringType,true)
和
org.apache.spark.sql.catalyst.errors.package $ TreeNodeException:未解析的属性
分别.
我有一个简单的JSON数据集,如下所示.如何查询所有parts.lock的id= 1.
JSON:
{
"id": 1,
"name": "A green door",
"price": 12.50,
"tags": ["home", "green"],
"parts" : [
{
"lock" : "One lock",
"key" : "single key"
},
{
"lock" : "2 lock",
"key" : "2 key"
}
]
}
Run Code Online (Sandbox Code Playgroud)
查询:
select id,name,price,parts.lockfrom product where id=1
Run Code Online (Sandbox Code Playgroud)
关键是如果我使用parts[0].lock它将返回如下一行:
{u'price': 12.5, u'id': 1, u'.lock': {u'lock': u'One lock', u'key': u'single key'}, u'name': u'A green door'}
Run Code Online (Sandbox Code Playgroud)
但我想返回所有locks的parts结构.它将返回多行,但这是我正在寻找的那一行.这种我想要完成的关系连接.
请在这件事上给予我帮助
如何在火花数据帧中投射结构数组?
让我通过一个例子来解释我想要做什么。我们将首先创建一个包含行数组和嵌套行的数据框。我的整数尚未在数据框中进行转换,它们被创建为字符串:
import org.apache.spark.sql._
import org.apache.spark.sql.types._
val rows1 = Seq(
Row("1", Row("a", "b"), "8.00", Seq(Row("1","2"), Row("12","22"))),
Row("2", Row("c", "d"), "9.00", Seq(Row("3","4"), Row("33","44")))
)
val rows1Rdd = spark.sparkContext.parallelize(rows1, 4)
val schema1 = StructType(
Seq(
StructField("id", StringType, true),
StructField("s1", StructType(
Seq(
StructField("x", StringType, true),
StructField("y", StringType, true)
)
), true),
StructField("d", StringType, true),
StructField("s2", ArrayType(StructType(
Seq(
StructField("u", StringType, true),
StructField("v", StringType, true)
)
)), true)
)
)
val df1 = spark.createDataFrame(rows1Rdd, schema1)
Run Code Online (Sandbox Code Playgroud)
这是创建的数据框的架构:
import org.apache.spark.sql._
import org.apache.spark.sql.types._
val rows1 = Seq(
Row("1", …Run Code Online (Sandbox Code Playgroud) 现在有JSON数据如下
{"Id":11,"data":[{"package":"com.browser1","activetime":60000},{"package":"com.browser6","activetime":1205000},{"package":"com.browser7","activetime":1205000}]}
{"Id":12,"data":[{"package":"com.browser1","activetime":60000},{"package":"com.browser6","activetime":1205000}]}
......
Run Code Online (Sandbox Code Playgroud)
此JSON是应用程序的激活时间,其目的是分析每个应用程序的总激活时间
我使用sparK SQL来解析JSON
斯卡拉
val sqlContext = sc.sqlContext
val behavior = sqlContext.read.json("behavior-json.log")
behavior.cache()
behavior.createOrReplaceTempView("behavior")
val appActiveTime = sqlContext.sql ("SELECT data FROM behavior") // SQL query
appActiveTime.show (100100) // print dataFrame
appActiveTime.rdd.foreach(println) // print RDD
Run Code Online (Sandbox Code Playgroud)
但是打印的dataFrame是这样的
.
+----------------------------------------------------------------------+
| data|
+----------------------------------------------------------------------+
| [[60000, com.browser1], [12870000, com.browser]]|
| [[60000, com.browser1], [120000, com.browser]]|
| [[60000, com.browser1], [120000, com.browser]]|
| [[60000, com.browser1], [1207000, com.browser]]|
| [[120000, com.browser]]|
| [[60000, com.browser1], [1204000, com.browser5]]|
| [[60000, com.browser1], [12075000, com.browser]]|
| [[60000, com.browser1], …Run Code Online (Sandbox Code Playgroud)