正如您所说,您在每一行中都有json 格式的问题
{"userId":"12345","vars":{"test_group":"group1","brand":"xband"},"modules":[{"id":"New"},{"id":"Default"},{"id":"BestValue"},{"id":"Rating"},{"id":"DeliveryMin"},{"id":"Distance"}]}
{"userId":"12345","vars":{"test_group":"group1","brand":"xband"},"modules":[{"id":"New"},{"id":"Default"},{"id":"BestValue"},{"id":"Rating"},{"id":"DeliveryMin"},{"id":"Distance"}]}
如果是这样,那么您可以使用sqlContext 的json api 将json 文件读取到dataframe,如下所示
val df = sqlContext.read.json("path to json file")
应该给你dataframe
+--------------------------------------------------------------------+------+--------------+
|modules |userId|vars |
+--------------------------------------------------------------------+------+--------------+
|[[New], [Default], [BestValue], [Rating], [DeliveryMin], [Distance]]|12345 |[xband,group1]|
|[[New], [Default], [BestValue], [Rating], [DeliveryMin], [Distance]]|12345 |[xband,group1]|
+--------------------------------------------------------------------+------+--------------+
和schema 是
root
|-- modules: array (nullable = true)
| |-- element: struct (containsNull = true)
| | |-- id: string (nullable = true)
|-- userId: string (nullable = true)
|-- vars: struct (nullable = true)
| |-- brand: string (nullable = true)
| |-- test_group: string (nullable = true)
最后一步将是 filter 仅将 modules.id 与 Default 作为值
val finaldf = df.withColumn("modules", explode($"modules.id"))
.filter($"modules" === "Default")
这应该给你
+-------+------+--------------+
|modules|userId|vars |
+-------+------+--------------+
|Default|12345 |[xband,group1]|
|Default|12345 |[xband,group1]|
+-------+------+--------------+
希望回答对你有帮助
更新
这将创建json
{"modules":"Default","userId":"12345","vars":{"brand":"xband","test_group":"group1"}}
{"modules":"Default","userId":"12345","vars":{"brand":"xband","test_group":"group1"}}
但如果你的要求是得到如下
{"modules":{"id":"Default"},"userId":"12345","vars":{"brand":"xband","test_group":"group1"}}
{"modules":{"id":"Default"},"userId":"12345","vars":{"brand":"xband","test_group":"group1"}}
你应该爆炸modules而不是modules.id
val finaldf = df.withColumn("modules", explode($"modules"))
.filter($"modules.id" === "Default")