【问题标题】:Spark java : Creating a new Dataset with a given schemaSpark java:使用给定模式创建新数据集
【发布时间】:2018-08-01 14:13:35
【问题描述】:

我有这段代码在 scala 中运行良好:

val schema = StructType(Array(
        StructField("field1", StringType, true),
        StructField("field2", TimestampType, true),
        StructField("field3", DoubleType, true),
        StructField("field4", StringType, true),
        StructField("field5", StringType, true)
    ))

val df = spark.read
    // some options
    .schema(schema)
    .load(myEndpoint)

我想在 Java 中做类似的事情。所以我的代码如下:

final StructType schema = new StructType(new StructField[] {
     new StructField("field1",  new StringType(), true,new Metadata()),
     new StructField("field2", new TimestampType(), true,new Metadata()),
     new StructField("field3", new StringType(), true,new Metadata()),
     new StructField("field4", new StringType(), true,new Metadata()),
     new StructField("field5", new StringType(), true,new Metadata())
});

Dataset<Row> df = spark.read()
    // some options
    .schema(schema)
    .load(myEndpoint);

但这给了我以下错误:

Exception in thread "main" scala.MatchError: org.apache.spark.sql.types.StringType@37c5b8e8 (of class org.apache.spark.sql.types.StringType)

我的架构似乎没有问题,所以我真的不知道问题出在哪里。

spark.read().load(myEndpoint).printSchema();
root
 |-- field5: string (nullable = true)
 |-- field2: timestamp (nullable = true)
 |-- field1: string (nullable = true)
 |-- field4: string (nullable = true)
 |-- field3: string (nullable = true)

schema.printTreeString();
root
 |-- field1: string (nullable = true)
 |-- field2: timestamp (nullable = true)
 |-- field3: string (nullable = true)
 |-- field4: string (nullable = true)
 |-- field5: string (nullable = true)

编辑:

这是一个数据样本:

spark.read().load(myEndpoint).show(false);
+---------------------------------------------------------------+-------------------+-------------+--------------+---------+
|field5                                                         |field2             |field1       |field4        |field3   |
+---------------------------------------------------------------+-------------------+-------------+--------------+---------+
|{"fieldA":"AAA","fieldB":"BBB","fieldC":"CCC","fieldD":"DDD"}  |2018-01-20 16:54:50|SOME_VALUE   |SOME_VALUE    |0.0      |
|{"fieldA":"AAA","fieldB":"BBB","fieldC":"CCC","fieldD":"DDD"}  |2018-01-20 16:58:50|SOME_VALUE   |SOME_VALUE    |50.0     |
|{"fieldA":"AAA","fieldB":"BBB","fieldC":"CCC","fieldD":"DDD"}  |2018-01-20 17:00:50|SOME_VALUE   |SOME_VALUE    |20.0     |
|{"fieldA":"AAA","fieldB":"BBB","fieldC":"CCC","fieldD":"DDD"}  |2018-01-20 18:04:50|SOME_VALUE   |SOME_VALUE    |10.0     |
 ...
+---------------------------------------------------------------+-------------------+-------------+--------------+---------+

【问题讨论】:

    标签: java scala apache-spark apache-spark-dataset


    【解决方案1】:

    使用 Datatypes 类中的静态方法和字段,而不是构造函数在 Spark 2.3.1 中为我工作:

        StructType schema = DataTypes.createStructType(new StructField[] {
                DataTypes.createStructField("field1",  DataTypes.StringType, true),
                DataTypes.createStructField("field2", DataTypes.TimestampType, true),
                DataTypes.createStructField("field3", DataTypes.StringType, true),
                DataTypes.createStructField("field4", DataTypes.StringType, true),
                DataTypes.createStructField("field5", DataTypes.StringType, true)
        });
    

    【讨论】:

    • 您好,感谢您的回答。您对 DoubleType 的看法是正确的,但这不是问题所在。使用构造函数而不是静态 DataTypes 函数。非常感谢!
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2018-12-06
    • 2019-05-04
    • 1970-01-01
    • 1970-01-01
    • 2021-01-28
    • 2021-03-03
    • 2019-10-16
    相关资源
    最近更新 更多