【问题标题】:filter count a spark dataframe过滤计数火花数据帧
【发布时间】:2020-08-09 22:57:30
【问题描述】:

我有两个如下的数据框,我正在从 MySQL 表中读取逻辑 DF

逻辑 DF:

slNo | filterCondtion |
-----------------------
1    | age > 100      |
2    | age > 50       |
3    | age > 10       |
4    | age > 20       |

InputDF - 从文件中读取:

age   | name           |
------------------------
11    | suraj          |
22    | surjeth        |
33    | sam            |
43    | ram            |

我想从逻辑数据框中应用过滤器语句并添加这些过滤器的计数

结果输出:

slNo | filterCondtion | count |
------------------------------
1    | age > 100      |   10  |
2    | age > 50       |   2   |
3    | age > 10       |   5   |
4    | age > 20       |   6   |
-------------------------------

我尝试过的代码:

val LogicDF = spark.read.format("jdbc").option("url", "jdbc:mysql://localhost:3306/testDB").option("driver", "com.mysql.jdbc.Driver").option("dbtable", "logic_table").option("user", "root").option("password", "password").load()

def filterCount(str: String): Long ={
     val counte = inputDF.where(str).count()
counte
}

val filterCountUDF = udf[Long, String](filterCount)

LogicDF.withColumn("count",filterCountUDF(col("filterCondtion")))

错误跟踪:

Caused by: org.apache.spark.SparkException: Failed to execute user defined function($anonfun$1: (string) => bigint)
  at org.apache.spark.sql.catalyst.expressions.GeneratedClass$GeneratedIteratorForCodegenStage1.processNext(Unknown Source)
  at org.apache.spark.sql.execution.BufferedRowIterator.hasNext(BufferedRowIterator.java:43)
  at org.apache.spark.sql.execution.WholeStageCodegenExec$$anonfun$11$$anon$1.hasNext(WholeStageCodegenExec.scala:619)
  at org.apache.spark.sql.execution.SparkPlan$$anonfun$2.apply(SparkPlan.scala:255)
  at org.apache.spark.sql.execution.SparkPlan$$anonfun$2.apply(SparkPlan.scala:247)
  at org.apache.spark.rdd.RDD$$anonfun$mapPartitionsInternal$1$$anonfun$apply$24.apply(RDD.scala:836)
  at org.apache.spark.rdd.RDD$$anonfun$mapPartitionsInternal$1$$anonfun$apply$24.apply(RDD.scala:836)
  at org.apache.spark.rdd.MapPartitionsRDD.compute(MapPartitionsRDD.scala:52)
  at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:324)
  at org.apache.spark.rdd.RDD.iterator(RDD.scala:288)
  at org.apache.spark.rdd.MapPartitionsRDD.compute(MapPartitionsRDD.scala:52)
  at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:324)
  at org.apache.spark.rdd.RDD.iterator(RDD.scala:288)
  at org.apache.spark.scheduler.ResultTask.runTask(ResultTask.scala:90)
  at org.apache.spark.scheduler.Task.run(Task.scala:121)
  at org.apache.spark.executor.Executor$TaskRunner$$anonfun$10.apply(Executor.scala:402)
  at org.apache.spark.util.Utils$.tryWithSafeFinally(Utils.scala:1360)
  at org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:408)
  at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149)
  at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624)
  at java.lang.Thread.run(Thread.java:748)
Caused by: java.lang.NullPointerException
  at org.apache.spark.sql.Dataset.where(Dataset.scala:1525)
  at filterCount(<console>:28)
  at $anonfun$1.apply(<console>:25)
  at $anonfun$1.apply(<console>:25)
  ... 21 more

任何选择也可以..!提前致谢。

【问题讨论】:

  • 您确定inputDF 不为空吗?
  • 是的。 InputDF 不为空。它正在工作,如果创建一个像这样val logicDF = Seq( (1, "age &gt; 18"), (2, "age &lt; 18") ).toDF("slno", "filterCondtion") 的 DF,但是当在从 MySQL 读取的 DF 上应用某些东西时..它会导致这个问题。
  • hm,不确定是否可以在 UDF 中使用另一个数据框。另见stackoverflow.com/questions/50123238/…stackoverflow.com/questions/41390572/…
  • 好的,有什么办法可以解决这个问题吗?
  • 请在使用 MySQL 读取数据后显示模式 DataFrames。

标签: scala apache-spark apache-spark-sql user-defined-functions


【解决方案1】:

没有 UDF 的解决方案

只要您的 logicDF 小到可以收集到驱动程序中,这将起作用。

步骤 1

将您的逻辑收集到Array[(Int, String)],如下:

val rules = logicDF.collect().map{ r: Row =>
  val slNo = r.getAs[Int](0)
  val condition = r.getAs[String](1)
  (slNo, condition)
}

第二步

使用条件值构建一个新列,将这些规则链接到 when Column。为此,请使用一些 scala 循环,例如:

val unused = when(lit(false), lit(false))
val filters: Column = rules.foldLeft(unused){
  case (acc: Column, (slNo: Int, cond: String)) =>
    acc.when(col("slNo") === slNo, expr(cond))
}

//You will get something like:
//when(col("slNo") === 1, expr("age > 10"))
//.when(col("slNo") === 2, expr("age > 20"))
//...

第三步

通过连接获取两个 DataFrame 的笛卡尔积,因此您可以将每条规则应用于数据中的每一行:

val joinDF = logicDF.join(inputDF, lit(true), "inner") //inner or whatever

第四步

使用以前的 Column 和条件过滤器进行过滤。

val withRulesDF = joinDF.filter(filters)

第 5 步

分组和计数:

val resultDF = withRulesDF
  .groupBy("slNo", "filterCondtion")
  .agg(count("*") as "count")

【讨论】:

    【解决方案2】:
    package spark
    
    import org.apache.spark.sql.{DataFrame, SparkSession}
    import org.apache.spark.sql.functions._
    
    object LogicFilterDataFrame extends App {
      val spark = SparkSession.builder()
        .master("local")
        .appName("DataFrame-example")
        .getOrCreate()
    
      import spark.implicits._
    
      case class LogicFilter(slNo: Int, filterCondition: String)
      case class Data(age: Int, name:String)
    
      val logicDF = Seq(
        LogicFilter(1, "age > 100"),
        LogicFilter(2, "age > 50"),
        LogicFilter(3, "age > 10"),
        LogicFilter(4, "age > 20")
      ).toDF()
    
      val dataDF = Seq(
        Data(11, "suraj"),
        Data(22, "surjeth"),
        Data(33, "sam"),
        Data(43, "ram")
      ).toDF()
    
      val logicCount = udf{s: String => {
        dataDF.filter(s).count()
        }}
      val resDF = logicDF.filter('filterCondition.like("%age%")).withColumn("count", logicCount('filterCondition))
      resDF.show(false)
    
    }
    

    【讨论】:

    • 谢谢。我尝试使用隐式 val logicDF = Seq( (1, "age &gt; 18"), (2, "age &lt; 18") ).toDF("slno", "filterCondtion") 在代码中创建 DF,它也适用于 mee。但是当我从 MySQL DB 中读取 logicDF 时,我遇到了问题。
    • LogicDF.withColumn("count",filterCountUDF(col("filterCondtion"))) 在您的代码中 -- filterCondition: string (nullable = true - 在您的架构中
    • 我也试过你的代码,如果我用隐式创建一个 logicDF,它就可以工作,但是如果我尝试从 MySQL 读取,我会得到那个空指针异常。仅供参考,即使在您的代码中 filterCondition 也是字符串,其可为空的 true
    • 另外,我们需要在字段filterCondition中使用MySQL检查数据。
    猜你喜欢
    • 2016-05-01
    • 2018-02-15
    • 2019-01-14
    • 2023-03-29
    • 2018-10-27
    • 2015-12-16
    • 1970-01-01
    • 2016-06-23
    • 1970-01-01
    相关资源
    最近更新 更多