【发布时间】:2020-01-19 10:03:05
【问题描述】:
我有一个字符串,其中包含需要进入我预期数据帧的.agg 函数的函数。
我的数据数据框看起来像
val client = Seq((1,"A","D",10),(2,"A","D",5),(3,"B","C",56),(5,"B","D",67)).toDF("ID","Categ","subCat","Amnt")
+---+-----+------+----+
| ID|Categ|subCat|Amnt|
+---+-----+------+----+
| 1| A| D| 10|
| 2| A| D| 5|
| 3| B| C| 56|
| 5| B| D| 67|
+---+-----+------+----+
所以我试图插入这个刺痛
val str= "s"$count(ID) as Total,$sum(Amnt) as amt""
我想实现这个作为输出
client.groupBy("Categ","subCat").agg(sum("Amnt") as "amt",count("ID") as "Total").show()
+-----+------+---+-----+
|Categ|subCat|amt|Total|
+-----+------+---+-----+
| B| C| 56| 1|
| A| D| 15| 2|
| B| D| 67| 1|
+-----+------+---+-----+
我试过了
client.groupBy("Categ","subCat").agg(s"$str").show()
出现错误
> error: overloaded method value agg with alternatives:
(expr: org.apache.spark.sql.Column,exprs: org.apache.spark.sql.Column*)org.apache.spark.sql.DataFrame
(exprs: java.util.Map[String,String])org.apache.spark.sql.DataFrame (表达式: scala.collection.immutable.Map[String,String])org.apache.spark.sql.DataFrame (aggExpr: (String, String),aggExprs: (String, String)*)org.apache.spark.sql.DataFrame 不能应用于(String)
我也试过expr
val str="sum(Amnt) as amt"
client.groupBy("Categ","subCat").agg(expr(str)).show()
this return the desired outcome
+-----+------+---+
|Categ|subCat|amt|
+-----+------+---+
| B| C| 56|
| A| D| 15|
| B| D| 67|
+-----+------+---+
但是当我再次尝试时
val str="sum(Amnt) as amt,count(ID) as ID_tot"
client.groupBy("Categ","subCat").agg(expr(str)).show()
org.apache.spark.sql.catalyst.parser.ParseException:
mismatched input ',' expecting <EOF>(line 1, pos 16)
【问题讨论】:
-
您正在尝试在此处混合使用 Spark SQL 和 Dataframe API。这是不可能的。如果您提到的要求严格要求字符串插值,那么您必须选择纯 Spark SQL 解决方案,即
select count(ID) as Total, sum(Amnt) as amt from client group by Categ ,subCat -
@Alexandros 我认为这是我的备用计划。
-
如您所见,here
agg有 3 个重载 1.agg(aggExpr: (String, String), aggExprs: (String, String)*): DataFrame2.agg(exprs: Map[String, String]): DataFrame3.agg(expr: Column, exprs: Column*): DataFrame它们都不接受字符串。所以你可以像上面显示的那样做df.groupBy("Categ","subCat").agg("Amnt" -> "sum", "ID" -> "count")或df.groupBy("Categ","subCat").agg(Map("Amnt" -> "sum", "ID" -> "count"))或SQL
标签: scala apache-spark apache-spark-sql string-interpolation