【问题标题】:Exploding a JSON array in a Spark Dataset在 Spark 数据集中分解 JSON 数组
【发布时间】:2017-04-11 22:19:26
【问题描述】:

我正在使用 Spark 2.1 和 Zeppelin 0.7 执行以下操作。 (这受到 Databricks 教程 (https://databricks.com/blog/2017/01/19/real-time-streaming-etl-structured-streaming-apache-spark-2-1.html) 的启发)

我创建了以下架构

val jsonSchema = new StructType()
.add("Records", ArrayType(new StructType()
    .add("Id", IntegerType)
    .add("eventDt", StringType)
    .add("appId", StringType)
    .add("userId", StringType)
    .add("eventName", StringType)
    .add("eventValues", StringType)
   )
  )

读取以下 json 'array' 文件,该文件位于我的 'inputPath' 目录中

{
"Records": [{
    "Id": 9550,
    "eventDt": "1491810477700",
    "appId": "dandb01",
    "userId": "985580",
    "eventName": "OG: HR: SELECT",
    "eventValues": "985087"
    },
    ... other records
]}

val rawRecords = spark.read.schema(jsonSchema).json(inputPath)

然后我想分解这些记录以获取各个事件

val events = rawRecords.select(explode($"Records").as("record"))

但是 rawRecords.show() 和 events.show() 都是空的。

知道我做错了什么吗?过去我知道我应该为此使用 JSONL,但 Databricks 教程建议最新版本的 spark 现在应该支持 json 数组。

【问题讨论】:

  • 实际上你的代码是有效的。这是你的 json 文件。 Spark 不喜欢格式化的 JSON。尝试格式化一个单行 json,它会起作用。

标签: json apache-spark dataset


【解决方案1】:

我做了以下事情:

  1. 我有一个包含以下数据的文件 foo.txt

{"记录":[{"Id":9550,"eventDt":"1491810477700","appId":"dandb01","userId":"985580","eventName":"OG:HR:SELECT ","eventValues":"985087"},{"Id":9550,"eventDt":"1491810477700","appId":"dandb01","userId":"985580","eventName":"OG:HR :SELECT","eventValues":"985087"},{"Id":9550,"eventDt":"1491810477700","appId":"dandb01","userId":"985580","eventName":"OG :HR:SELECT","eventValues":"985087"},{"Id":9550,"eventDt":"1491810477700","appId":"dandb01","userId":"985580","eventName": "OG:HR:SELECT","eventValues":"985087"}]} {"记录":[{"Id":9550,"eventDt":"1491810477700","appId":"dandb01","userId":"985580","eventName":"OG:HR:SELECT"," eventValues":"985087"},{"Id":9550,"eventDt":"1491810477700","appId":"dandb01","userId":"985580","eventName":"OG:HR:SELECT" ,"eventValues":"985087"},{"Id":9550,"eventDt":"1491810477700","appId":"dandb01","userId":"985580","eventName":"OG:HR: SELECT","eventValues":"985087"},{"Id":9550,"eventDt":"1491810477700","appId":"dandb01","userId":"985580","eventName":"OG: HR:SELECT","eventValues":"985087"}]}

  1. 我有以下代码

    导入 sqlContext.implicits._ 导入 org.apache.spark.sql.functions._

    val df = sqlContext.read.json("foo.txt") df.printSchema()
    df.select(explode($"Records").as("record")).show

  2. 我得到以下输出

root |-- 记录:数组(可为空=真)| |-- 元素:结构 (包含空=真)| | |-- id: long (nullable = true) |
| |-- appId: 字符串 (可为空 = true) | | |-- eventDt: 字符串(可为空=真)| | |-- 事件名称:字符串(可为空 = 真的)| | |-- eventValues: 字符串 (可为空 = true) | |
|-- 用户ID:字符串(可为空=真)

+--------------------+
|              record|
+--------------------+
|[9550,dandb01,149...|
|[9550,dandb01,149...|
|[9550,dandb01,149...|
|[9550,dandb01,149...|
|[9550,dandb01,149...|
|[9550,dandb01,149...|
|[9550,dandb01,149...|
|[9550,dandb01,149...|
+--------------------+

【讨论】:

    猜你喜欢
    • 2016-05-06
    • 2019-02-21
    • 1970-01-01
    • 2017-03-24
    • 1970-01-01
    • 2019-11-01
    • 2017-11-10
    • 2021-01-14
    • 2018-01-19
    相关资源
    最近更新 更多