【发布时间】:2020-12-06 13:08:16
【问题描述】:
我在 java8 中使用 spark-sql-2.4.1v。我有一个场景,如果列出现在给定的数据框列列表中,我需要执行某些操作
我有如下示例数据框,数据框的列会根据在数据库表上执行的外部查询而有所不同。
val data = List(
("20", "score", "school", "2018-03-31", 14 , 12 , 20),
("21", "score", "school", "2018-03-31", 13 , 13 , 21),
("22", "rate", "school", "2018-03-31", 11 , 14, 22),
("21", "rate", "school", "2018-03-31", 13 , 12, 23)
)
val df = data.toDF("id", "code", "entity", "date", "column1", "column2" ,"column3"..."columnN")
如上图所示,数据框“数据”列不是固定的,并且会有所不同,并且会有“column1”、“column2”、“column3”...“columnN”...
所以取决于列的可用性,我需要执行一些操作 同样,我尝试使用“when-clause”,当存在列时,我必须对指定列执行某些操作,否则继续下一个操作..
我正在尝试以下两种使用“when-cluase”的方法
第一路:
Dataset<Row> resultDs = df.withColumn("column1_avg",
when( df.schema().fieldNames().contains(col("column1")) , avg(col("column1"))))
)
第二种方式:
Dataset<Row> resultDs = df.withColumn("column2_sum",
when( df.columns().contains(col("column2")) , sum(col("column1"))))
)
错误:
无法在数组类型 String[] 上调用 contains(Column)
那么如何使用 java8 代码来处理这种情况呢?
【问题讨论】:
-
请显示预期输出
-
@thebluephantom 预期输出是动态的,取决于列...即,如果“column1”存在,那么我将对 column1 进行平均,如果 column2 存在,我将对列等进行求和......棘手这里的部分是如果列没有“column1”操作不应该失败因此我需要检查列是否存在
-
函数
sum和avg是聚合函数。他们通常会返回一行。因此,您必须编辑您的问题,并展示您期望的输出示例。你可以做两个例子——一个在 column1 存在时,一个在它不存在时。显示输入 ds 和输出 ds。
标签: dataframe apache-spark apache-spark-sql spark-streaming