【问题标题】:Group by inside otherwise clause, spark java在 else 子句中分组,火花 java
【发布时间】:2021-11-21 10:34:37
【问题描述】:

我在 SparkJava(IntelliJ 应用程序)中有这个过程,我遇到了一个我不知道如何解决的问题。首先我声明数据集:

private static final String CONTRA1 = "contra1";

query = "select contra1, ..., eadfinal, , ..., data_date" + FROM + dbSchema + TBLNAME " + WHERE fech = '" + fechjmCto2 + "' AND s1emp=49";
        Dataset<Row> jmCto2 = sql.sql(query);

然后我进行计算,我分析一些字段以分配一些文字值。我的问题出在聚合函数中:

Dataset<Row> contrCapOk1 = contrCapOk.join(jmCto2,
        contrCapOk.col(CONTRA1).equalTo(jmCto2.col(CONTRA1)),LEFT)
        .select(contrCapOk.col("*"),
        jmCto2.col("ind"),
 
functions.when(jmCto2.col(CONTRA1).isNull(),functions.lit(NUEVES))
     .when(jmCto2.col("ind").equalTo("N"),functions.lit(UNOS))
     .otherwise(jmCto2.groupBy(CONTRA1).agg(functions.sum(jmCto2.col("eadfinal")))).as("EAD"),

我想要的是在其他部分中求和。但是当我执行集群时,在日志中给我这个消息。

User class threw exception: java.lang.RuntimeException: Unsupported literal type class org.apache.spark.sql.Dataset [contra1: int, sum(eadfinal): decimal(33,6)] 

在第 211 行,否则在行。

你知道问题可能是什么吗?

谢谢。

【问题讨论】:

    标签: java apache-spark intellij-14 spark-java


    【解决方案1】:

    您不能在列子句中使用groupBy 和聚合函数。要做你想做的事,你必须使用window

    对于您的情况,您可以定义以下窗口:

    import org.apache.spark.sql.expressions.Window;
    import org.apache.spark.sql.expressions.WindowSpec;
    
    ...
    
    WindowSpec window = Window
      .partitionBy(CONTRA1)
      .rangeBetween(Window.unboundedPreceding(), Window.unboundedFollowing());
    

    在哪里

    • partitionBy 相当于 groupBy 用于聚合
    • rangeBetween 确定分区的哪些行将被聚合函数使用,这里我们取所有行

    然后你在调用你的聚合函数时使用这个窗口,如下:

    import org.apache.spark.sql.functions;
    
    ...
    
    Dataset<Row> contrCapOk1 = contrCapOk.join(
        jmCto2,
        contrCapOk.col(CONTRA1).equalTo(jmCto2.col(CONTRA1)),
        LEFT
      )
      .select(
        contrCapOk.col("*"),
        jmCto2.col("ind"),
        functions.when(jmCto2.col(CONTRA1).isNull(), functions.lit(NUEVES))
         .when(jmCto2.col("ind").equalTo("N"), functions.lit(UNOS))
         .otherwise(functions.sum(jmCto2.col("eadfinal")).over(window))
         .as("EAD")
      )
    

    【讨论】:

    • 感谢 Vicent,但它无法正常工作。我忘了在代码中提到上层,在数据集中,我有这些行“ jmCto2 = jmCto2.withColumn("rank", functions.rank().over(rank0)); jmCto2 = jmCto2.selectht("*" ).where(jmCto2.col("rank").equalTo(1)); ".我不确定,但我认为这是在避免赚钱。我想我将不得不制作一个具有相同过滤器但没有 rank 子句的新数据集。是这样吗?
    • “它不能正常工作”是什么意思?您的应用程序是否抛出错误?它是否运行顺利,但您没有得到预期的结果?关于您的rank 用法,我看不出有任何问题。你如何定义变量rank0
    • 你好,没有错误信息。代码抛出一个文件的结果,没有总和。这是代码的那一部分: WindowSpec rank0 = Window.partitionBy(CONTRA1,FEOPERAC).orderBy(jmCto2.col(DATA_DATE_PART).desc()); jmCto2 = jmCto2.withColumn("rank", functions.rank().over(rank0)); jmCto2 = jmCto2.select("*").where(jmCto2.col("rank").equalTo(1)); jmCto2 = jmCto2.drop("排名");
    • 我不认为您的问题来自rank 部分代码,它似乎完全合法。但是,我仍然不明白您对选择的确切期望。您可以添加一些输入数据帧的样本(您应用选择的数据帧)以及您对问题中的输出数据帧的期望是什么?它会帮助我得到你想做的事情。
    • 你好,我一直在检查我的代码我做了一些调整。我应用了你的部分代码,结果非常棒,这就是我想要的。再次感谢您,如果我的解释不是很准确,请原谅。- 谢谢。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2015-06-05
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2017-02-21
    • 1970-01-01
    相关资源
    最近更新 更多