【问题标题】:Spark Dataframe datatype as StringSpark Dataframe 数据类型为字符串
【发布时间】:2018-12-24 02:45:53
【问题描述】:

我正在尝试通过将描述写入 SQL 查询来验证 DataFrame 的数据类型,但每次我将日期时间作为字符串获取。

1.首先我尝试了以下代码:

    SparkSession sparkSession=new SparkSession.Builder().getOrCreate();
        Dataset<Row> df=sparkSession.read().option("header","true").option("inferschema","true").format("csv").load("/user/data/*_ecs.csv");

        try {
    df.createTempView("data");
    Dataset<Row> sqlDf=sparkSession.sql("Describe data");
    sqlDf.show(300,false);

    Output:
    +-----------------+---------+-------+
    |col_name         |data_type|comment|
    +-----------------+---------+-------+
    |id               |int      |null   |
    |symbol           |string   |null   |
    |datetime         |string   |null   |
    |side             |string   |null   |
    |orderQty         |int      |null   |
    |price            |double   |null   | 
    +-----------------+---------+-------+
  1. 我也尝试自定义架构,但在这种情况下,当我执行除描述表之外的任何查询时,我会遇到异常:

    SparkSession sparkSession=new SparkSession.Builder().getOrCreate(); Dataset<Row>df=sparkSession.read().option("header","true").schema(customeSchema).format("csv").load("/use/data/*_ecs.csv");
     try {
                df.createTempView("trade_data");
        Dataset<Row> sqlDf=sparkSession.sql("Describe trade_data");
        sqlDf.show(300,false);
    
    Output:
    +--------+---------+-------+
    |col_name|data_type|comment|
    +--------+---------+-------+
    |datetime|timestamp|null   |
    |price   |double   |null   |
    |orderQty|double   |null   |
    +--------+---------+-------+
    

但是如果我尝试任何查询然后得到以下执行:

Dataset<Row> sqlDf=sparkSession.sql("select DATE(datetime),avg(price),avg(orderQty) from data group by datetime");


java.lang.IllegalArgumentException
        at java.sql.Date.valueOf(Date.java:143)
        at org.apache.spark.sql.catalyst.util.DateTimeUtils$.stringToTime(DateTimeUtils.scala:137)

如何解决?

【问题讨论】:

  • 我使用的自定义模式:StructType customeSchema=new StructType(new StructField[] { new StructField("datetime",DataTypes.TimestampType,true,Metadata.empty()), new StructField("price ",DataTypes.DoubleType,true,Metadata.empty()), new StructField("orderQty",DataTypes.DoubleType,true,Metadata.empty())});
  • valueOf javadoc 清楚地说明了何时抛出 IllegalArgumentException。请验证您的数据。 docs.oracle.com/javase/8/docs/api/java/sql/…
  • @DanW:你能说出为什么 inferschema 不起作用吗??

标签: java apache-spark dataframe apache-spark-sql


【解决方案1】:
  1. 为什么 Inferschema 不起作用??

  2. 如果您不想提交自己的架构,一种方法是:

    Dataset<Row> df = sparkSession.read().format("csv").option("header","true").option("inferschema", "true").load("example.csv");
    
    df.printSchema();  // check output - 1
    df.createOrReplaceTempView("df");
    Dataset<Row> df1 = sparkSession.sql("select * , Date(datetime) as datetime_d from df").drop("datetime");
    df1.printSchema();  // check output - 2
    
    ====================================
    
    output - 1:
    root
     |-- id: integer (nullable = true)
     |-- symbol: string (nullable = true)
     |-- datetime: string (nullable = true)
     |-- side: string (nullable = true)
     |-- orderQty: integer (nullable = true)
     |-- price: double (nullable = true)
    
    output - 2:
    root
     |-- id: integer (nullable = true)
     |-- symbol: string (nullable = true)
     |-- side: string (nullable = true)
     |-- orderQty: integer (nullable = true)
     |-- price: double (nullable = true)
     |-- datetime_d: date (nullable = true)
    

    如果要转换的字段数不高,我会选择这种方法。

  3. 如果您想提交自己的架构:

    List<org.apache.spark.sql.types.StructField> fields = new ArrayList<>();
    fields.add(DataTypes.createStructField("datetime", DataTypes.TimestampType, true));
    fields.add(DataTypes.createStructField("price",DataTypes.DoubleType,true));
    fields.add(DataTypes.createStructField("orderQty",DataTypes.DoubleType,true));
    StructType schema = DataTypes.createStructType(fields);
    Dataset<Row> df = sparkSession.read().format("csv").option("header", "true").schema(schema).load("example.csv");
    
    df.printSchema(); // output - 1
    
    df.createOrReplaceTempView("df");
    Dataset<Row> df1 = sparkSession.sql("select * , Date(datetime) as datetime_d from df").drop("datetime");
    df1.printSchema(); // output - 2
    
    ======================================
    output - 1:
    root
     |-- datetime: timestamp (nullable = true)
     |-- price: double (nullable = true)
     |-- orderQty: double (nullable = true)
    
    output - 2:
    root
     |-- price: double (nullable = true)
     |-- orderQty: double (nullable = true)
     |-- datetime_d: date (nullable = true)
    

    由于它再次将列从时间戳转换为日期,我看不到这种方法有太多用途。但还是把它放在这里供你以后使用。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2022-01-05
    • 1970-01-01
    • 2020-01-02
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2018-10-01
    相关资源
    最近更新 更多