【发布时间】: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 > 18"), (2, "age < 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