【问题标题】:How to explode column with multiple records into multiple Columns in Spark如何在 Spark 中将具有多条记录的列分解为多个列
【发布时间】:2022-07-31 18:13:56
【问题描述】:

我是使用 Spark 和 Scala 的新手,希望就这种情况获得一些帮助: 这是我当前的架构。

 |-- _id: struct (nullable = true)
 |    |-- oid: string (nullable = true)
 |-- date: timestamp (nullable = true)
 |-- horizon: double (nullable = true)
 |-- risk_table: array (nullable = true)
 |    |-- element: struct (containsNull = true)
 |    |    |-- index: string (nullable = true)
 |    |    |-- risk_buy: double (nullable = true)
 |    |    |-- reward_buy: double (nullable = true)
 |    |    |-- risk_sell: double (nullable = true)
 |    |    |-- reward_sell: double (nullable = true)
 |-- symbol_id: string (nullable = true)

以下是数据外观的示例:

+--------------------+
|          risk_table|
+--------------------+
|[{count, 201.0, 2...|
|[{count, 219.0, 2...|
|[{count, 119.0, 1...|
|[{count, 217.0, 2...|
|[{count, 17.0, 17...|
|[{count, 189.0, 1...|
|[{count, 105.0, 1...|
|[{count, 188.0, 1...|
|[{count, 111.0, 1...|
|[{count, 276.0, 2...|
|[{count, 70.0, 70...|
|[{count, 121.0, 1...|
|[{count, 133.0, 1...|
|[{count, 116.0, 1...|
|[{count, 70.0, 70...|
|[{count, 193.0, 1...|
|[{count, 131.0, 1...|
|[{count, 93.0, 93...|
|[{count, 84.0, 84...|
|[{count, 114.0, 1...|
+--------------------+

我想将 risk_table 列值分解为多个列,通常有 4 个嵌套文档/字典,其中索引名称发生变化,因此预期的输出看起来像这样

+-----------+------+---------+------------------+--------------------+-----+---------------------+
| symbol_id | date | index_0 | risk_buy_index_0 | reward_buy_index_0 | ... | reward_sell_index_3 |
+-----------+------+---------+------------------+--------------------+-----+---------------------+
| APPL      | xxxx | 0       | 0                | 0                  | ... | 0                   |
+-----------+------+---------+------------------+--------------------+-----+---------------------+
| APPL      | xxxx | 0       | 0                | 0                  | ... | 0                   |
+-----------+------+---------+------------------+--------------------+-----+---------------------+
| APPL      | xxxx | 0       | 0                | 0                  | ... | 0                   |
+-----------+------+---------+------------------+--------------------+-----+---------------------+


我找到了一些关于如何只分解一个文档/字典而不是嵌套的信息,如果有人能提供帮助,我将不胜感激。

【问题讨论】:

    标签: json dataframe scala apache-spark apache-spark-sql


    【解决方案1】:

    假设您的数据集是main。首先,我们必须炸开risk_table 的内容,因为如果我们不这样做,我们会得到数组作为列的值,这是我们不喜欢的,所以:

    df1 = df1.withColumn("explode", explode(col("risk_table")))
    

    现在,explode 列每行有一个对象;有很多方法可以从对象创建列,但我喜欢使用 selectExpr:

    .selectExpr("id", "symbol_id", // or whatever other field you like
      "explode.index as index_0",  // then target the key with dot operator
      "explode.risk_buy as risk_buy_index_0",
      "explode.reward_buy as reward_buy_index_0"
      // add your other wanted values
    )
    

    虚拟输入:

    +--------------------------+---+---------+
    |risk_table                |id |symbol_id|
    +--------------------------+---+---------+
    |[{1, 0.25, 0.3, 0.1, 0.3}]|1  |1        |
    +--------------------------+---+---------+
    

    最终输出:

    +---+---------+-------+----------------+------------------+
    | id|symbol_id|index_0|risk_buy_index_0|reward_buy_index_0|
    +---+---------+-------+----------------+------------------+
    |  1|        1|      1|            0.25|               0.3|
    +---+---------+-------+----------------+------------------+
    

    【讨论】:

      猜你喜欢
      • 2022-12-05
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2022-12-05
      • 1970-01-01
      • 2021-02-24
      • 1970-01-01
      相关资源
      最近更新 更多