【问题标题】:Convert JSON Data into DataFrame Apache Spark将 JSON 数据转换为 DataFrame Apache Spark
【发布时间】:2021-03-26 13:53:33
【问题描述】:

我想将我的 json 数据转换为数据框,以使其更易于管理。

将数据导入java的命令

Dataset<Row> df = spark.read()
   .option("multiline",true)
   .option("mode", "PERMISSIVE")
   .json("hdfs://hd-master:9820/houseInformation.txt");

数据

{
 "House1": {
   "House_Id": "1",
   "Cover": "1.000",
   "HouseType": "bungalow",
   "Facing": "South",
   "Region": "YVR",
   "Ru": "1",
   "HVAC": [
     "FAGF",
     "FPG",
     "HP"
   ]
 },
 "House2" : {...},
 "House3" : {...},
}

如果可能的话,我想删除键“House1”,然后将其余数据转换为 df。如果没有,那也没关系。

但理想情况下,这是我想要的输出

HouseName  House_Id  Cover   HouseType Facing Region Ru HVAC
 House1       1     1.000     bungalow  South   YVR   1  []   

架构

root
 |-- House1: struct (nullable = true)
 |    |-- Cover: string (nullable = true)
 |    |-- Facing: string (nullable = true)
 |    |-- HVAC: array (nullable = true)
 |    |    |-- element: string (containsNull = true)
 |    |-- HouseType: string (nullable = true)
 |    |-- House_Id: string (nullable = true)
 |    |-- Region: string (nullable = true)
 |    |-- Ru: string (nullable = true)

简单打印

df.select(functions.col("House1")).show(false);

返回这个

+----------------------------------------------------+
|House1                                              |
+----------------------------------------------------+
|[1.000, South, [FAGF, FPG, HP], bungalow, 1, YVR, 1]|
+----------------------------------------------------+

【问题讨论】:

  • house 列是 struct 类型,要查看 struct 的所有列 Try - df.select(functions.col("House1.*")).show(false);

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


【解决方案1】:

您的输入 json 文件应为 jsonlines 格式 (https://jsonlines.org/),而这是单个 json 文档。

如果所有这些“house”元素都具有相同的键并且文件是 jsonlines 格式,那么它们将自动成为 Dataframe 中的单独行。 但是这里的键是不同的,即 House1、House2 等,因此它被视为单个记录。 有任意数量的“房屋”,但每个都有不同的键,因此它们被视为不同的列。

如果您编写复杂的转换来实现您想要的结果,那将是没有效率的。单个 json 记录不会被 spark 并行化。将有一个执行者完成所有工作,并且您的数据框将具有任意数量的列。对于这个用例,你最好不要使用 spark。

因此,IMO 解决此问题的更明智方法是修复您的源文档。这可能意味着可能会预处理您的源文件以将其转换为jsonlines。

【讨论】:

  • 尽管所有其他答案都给了我想要的答案,但我还是选择了这个。这是最有效的。我将文件修改为适当的格式,并将更改通知我的教授。
  • 另一个可能对某人有所帮助的提示。有时您可能会得到一个带有 json 数组的输入 json 文件。通过删除文件中的第一个和最后一个字符,即删除文件开头的[ 和文件末尾的],可以轻松地将其转换为jsonlines
【解决方案2】:

由于House1 列的类型为struct,因此要从struct 类型中提取所有列,请使用House1.*。

试试下面的代码。

df
.select(functions.col("House1.*"))
.show(false)

【讨论】:

    【解决方案3】:

    使用 Scala

    1.读取众议院 JSON 数据。请注意,我在单行中给出一个 json 输入行

    { "House1": { "House_Id": "1", "Cover": "1.000", "HouseType": "bungalow", "Facing": "South", "Region": "YVR", "Ru": "1", "HVAC": [ "FAGF", "FPG", "HP" ] } }
    { "House2": { "House_Id": "2", "Cover": "1.000", "HouseType": "bungalow", "Facing": "North", "Region": "YVR", "Ru": "1", "HVAC": [ "FAGF", "FPG", "HP" ] } }
    

    代码

    val houseDS = spark.read.json("<JSON_FILE_PATH>");
    houseDS.printSchema
    root
     |-- House1: struct (nullable = true)
     |    |-- Cover: string (nullable = true)
     |    |-- Facing: string (nullable = true)
     |    |-- HVAC: array (nullable = true)
     |    |    |-- element: string (containsNull = true)
     |    |-- HouseType: string (nullable = true)
     |    |-- House_Id: string (nullable = true)
     |    |-- Region: string (nullable = true)
     |    |-- Ru: string (nullable = true)
     |-- House2: struct (nullable = true)
     |    |-- Cover: string (nullable = true)
     |    |-- Facing: string (nullable = true)
     |    |-- HVAC: array (nullable = true)
     |    |    |-- element: string (containsNull = true)
     |    |-- HouseType: string (nullable = true)
     |    |-- House_Id: string (nullable = true)
     |    |-- Region: string (nullable = true)
     |    |-- Ru: string (nullable = true)
    
    houseDS.show(false)
    +----------------------------------------------------+----------------------------------------------------+
    |House1                                              |House2                                              |
    +----------------------------------------------------+----------------------------------------------------+
    |[1.000, South, [FAGF, FPG, HP], bungalow, 1, YVR, 1]|null                                                |
    |null                                                |[1.000, North, [FAGF, FPG, HP], bungalow, 2, YVR, 1]|
    +----------------------------------------------------+----------------------------------------------------+
    

    2。我们使用stack() 函数将多列分隔为行。这是stack 函数语法:stack(n, expr1, ..., exprk) - 将expr1、...、exprk 分隔成n 行。

    val houseDS2 = houseDS.select(expr("stack(2,House1, 'House1', House2, 'House2') as (house,HouseName)")).na.drop
    houseDS2.printSchema
    root
     |-- house: struct (nullable = true)
     |    |-- Cover: string (nullable = true)
     |    |-- Facing: string (nullable = true)
     |    |-- HVAC: array (nullable = true)
     |    |    |-- element: string (containsNull = true)
     |    |-- HouseType: string (nullable = true)
     |    |-- House_Id: string (nullable = true)
     |    |-- Region: string (nullable = true)
     |    |-- Ru: string (nullable = true)
     |-- HouseName: string (nullable = true)
    

    3.然后从houseDS2DataSet 上面选择所有需要的列

    val finalHouseDS = houseDS2.select("HouseName","house.House_Id","house.Cover","house.HouseType","house.Facing","house.Region","house.Ru","house.HVAC")
    finalHouseDS.show(false)
    

    您的预期输出

    +---------+--------+-----+---------+------+------+---+---------------+
    |HouseName|House_Id|Cover|HouseType|Facing|Region|Ru |HVAC           |
    +---------+--------+-----+---------+------+------+---+---------------+
    |House1   |1       |1.000|bungalow |South |YVR   |1  |[FAGF, FPG, HP]|
    |House2   |2       |1.000|bungalow |North |YVR   |1  |[FAGF, FPG, HP]|
    +---------+--------+-----+---------+------+------+---+---------------+
    

    您可以在 Java 中进行类似的实现。如果您遇到较大数据集的任何性能问题,请告诉我。

    使用 Java

    import org.apache.spark.sql.Dataset;
    import org.apache.spark.sql.Row;
    import org.apache.spark.sql.SparkSession;
    import org.apache.spark.sql.types.StructField;
    import org.apache.spark.sql.functions.*;
    
    public class ParseJson {
        public static void main(String[] args) {
            System.setProperty("hadoop.home.dir", "D:\\Software\\Hadoop");
    
            SparkSession spark = SparkSession
                    .builder()
                    .appName("Testing")
                    .master("local[*]")
                    .getOrCreate();
            // Read json data
    
            Dataset<Row> houseDS = spark.read().json("<JSON_FILE_PATH>");
            houseDS.printSchema();
            Dataset<Row> houseDS2 = houseDS.selectExpr("stack(2,House1, 'House1', House2, 'House2') as (house,HouseName)").na().drop();
            houseDS2.printSchema();
            Dataset<Row> finalHouseDS = houseDS2.select("HouseName","house.House_Id","house.Cover","house.HouseType","house.Facing","house.Region","house.Ru","house.HVAC");
            finalHouseDS.show(false);
    
        }
    }
    
    

    【讨论】:

      【解决方案4】:

      您可以使用星号表示法选择结构的所有元素并展开为列:

      Dataset<Row> df2 = df.select("House1.*")
      
      df2.show(false)
      +-----+------+---------------+---------+--------+------+---+
      |Cover|Facing|HVAC           |HouseType|House_Id|Region|Ru |
      +-----+------+---------------+---------+--------+------+---+
      |1.000|South |[FAGF, FPG, HP]|bungalow |1       |YVR   |1  |
      +-----+------+---------------+---------+--------+------+---+
      

      【讨论】:

        猜你喜欢
        • 2021-12-23
        • 2016-05-16
        • 1970-01-01
        • 2017-03-17
        • 2023-03-31
        • 1970-01-01
        • 2016-04-11
        • 2016-04-24
        • 2021-07-04
        相关资源
        最近更新 更多