【问题标题】:Applying when condition only when column exists in the dataframe仅当数据框中存在列时才应用条件
【发布时间】: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”操作不应该失败因此我需要检查列是否存在
  • 函数sumavg聚合函数。他们通常会返回一行。因此,您必须编辑您的问题,并展示您期望的输出示例。你可以做两个例子——一个在 column1 存在时,一个在它不存在时。显示输入 ds 和输出 ds。

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


【解决方案1】:

您可以创建一个包含所有列名的列。然后您可以检查该列是否存在并处理它是否可用-

 df.withColumn("columns_available", array(df.columns.map(lit): _*))
      .withColumn("column1_org",
      when( array_contains(col("columns_available"),"column1") , col("column1")))
      .withColumn("x",
        when( array_contains(col("columns_available"),"column4") , col("column1")))
      .withColumn("column2_new",
        when( array_contains(col("columns_available"),"column2") , sqrt("column2")))
      .show(false)

【讨论】:

  • 非常感谢......但是如果没有指定的列存在如下所示的一个小疑问 df.withColumn("columns_available", array(df.columns.map(lit): _*) ) .withColumn("column1_org", when( array_contains(col("columns_available"),"columnN") , col("column2"))) .show(false) ,它不应该显示该列,即“column1_org”如何实现它?
  • 您无法使用withColumn 实现。如果该列不可用,为什么要首先添加withColumn。你能具体说明这个非常奇怪的要求的动机吗?
  • 我有基于某些可用列的要求,我需要执行某些操作,如果列不可用,则该操作未执行,因为操作特定于某些列。对于某些列,如果存在无效数据,它也可能为空......但如果我在这里添加空,那么它可能会被视为无效数据,但这里不是这种情况......
  • df.columns.map(lit): _* this 在这里做什么?当我在 java array(Arrays.asList(df.columns()).stream().map(s -> new Column(s)).toArray(Column[]::new)) 中执行它时,它实际获取列值而不是列
  • 你能告诉我这个广播变量访问有什么问题吗? stackoverflow.com/questions/64003697/…
猜你喜欢
  • 1970-01-01
  • 2019-09-16
  • 2016-04-11
  • 2021-02-18
  • 2021-02-24
  • 2013-01-10
  • 2014-02-22
  • 2020-11-27
  • 1970-01-01
相关资源
最近更新 更多