Spark流式窗口操作中指定多列groupBy报错,求问题排查
Spark groupBy 混合列名与窗口操作报错的解决办法
嘿,这个问题我之前也碰到过!咱们先拆解下报错的核心原因:
Spark的groupBy方法只有两种合法的调用方式:
- 第一种:传入字符串类型的列名(单个或多个,比如
groupBy("col1", "col2")) - 第二种:传入Column类型的对象(单个或多个,比如
groupBy(col("col1"), window(...)))
你现在的代码里,把List[String](列名列表)和Column(窗口操作)直接一起传给groupBy,这两种重载都不买账——因为前者只接受字符串参数,后者只接受Column参数,混合类型自然就触发错误了。
解决办法:把所有分组项统一转换成Column类型
只需要把你的字符串列名列表转换成Column对象,再和窗口列合并,一起传给groupBy就行。修正后的代码如下:
val groupCols = List("SINR_Distribution","NE_VERSION","NE_ID","NE_NAME","cNum","EarfcnDl","datetime","circle") // 把字符串列名转为Column对象,再追加窗口列 val groupByColumns = groupCols.map(col) :+ window($"EVENT_TIME", "60 minutes") // 使用:_*把序列拆解成可变参数传给groupBy val aggDFrame = dframe.groupBy(groupByColumns:_*).agg(Rule_Agg)
这里的:_*是Scala里的语法糖,用来把一个序列(比如List[Column])转换成方法需要的可变参数(Column*),这样groupBy就能正确识别所有分组项了。
再补一句为什么原来的写法不行
你之前的调用groupBy(groupCols, window(...))相当于给groupBy传了两个参数:第一个是List[String],第二个是Column。但groupBy的两个重载都不接受这种参数组合:
- 第一个重载要求参数是
String, String*,也就是第一个是字符串,后面跟着任意多个字符串 - 第二个重载要求参数是
Column*,也就是任意多个Column对象
类型完全不匹配,自然就抛出了你看到的错误。
内容的提问来源于stack exchange,提问作者vivekdesai
相关产品推荐
相关产品推荐

