过滤DataFrame最有效的方法是什么

Mar*_*aci 31 apache-spark apache-spark-sql

...通过检查列的值是否在a中seq.
也许我没有解释得很好,我基本上希望这(使用常规的SQL表达出来)DF_Column IN seq

首先,我使用a broadcast var(我放置seq),UDF(检查完成)和registerTempTable.
问题是,我没有测试它,因为我遇到了一个已知的bug,显然只有使用时,会出现registerTempTableScalaIDE.

我最终创建了一个新DataFrameseq并与它进行内连接(交集),但我怀疑这是完成任务的最高效的方式.

谢谢

编辑:(回应@YijieShen):
如何filter根据一个DataFrame列的元素是否在另一个DF的列(如SQL select * from A where login in (select username from B))中?

例如:第一个DF:

login      count
login1     192  
login2     146  
login3     72   
Run Code Online (Sandbox Code Playgroud)

第二个DF:

username
login2
login3
login4
Run Code Online (Sandbox Code Playgroud)

结果:

login      count
login2     146  
login3     72   
Run Code Online (Sandbox Code Playgroud)

尝试:
EDIT-2:我认为,现在修复了这个bug,这些应该可行.结束编辑-2

ordered.select("login").filter($"login".contains(empLogins("username")))
Run Code Online (Sandbox Code Playgroud)

ordered.select("login").filter($"login" in empLogins("username"))
Run Code Online (Sandbox Code Playgroud)

两者Exception in thread "main" org.apache.spark.sql.AnalysisException分别抛出:

resolved attribute(s) username#10 missing from login#8 in operator 
!Filter Contains(login#8, username#10);
Run Code Online (Sandbox Code Playgroud)

resolved attribute(s) username#10 missing from login#8 in operator 
!Filter login#8 IN (username#10);
Run Code Online (Sandbox Code Playgroud)

yjs*_*hen 16

我的代码(遵循第一种方法的描述)Spark 1.4.0-SNAPSHOT在这两种配置下正常运行:

  • Intellij IDEA's test
  • Spark Standalone cluster 有8个节点(1个主人,7个工人)

请检查是否存在任何差异

val bc = sc.broadcast(Array[String]("login3", "login4"))
val x = Array(("login1", 192), ("login2", 146), ("login3", 72))
val xdf = sqlContext.createDataFrame(x).toDF("name", "cnt")

val func: (String => Boolean) = (arg: String) => bc.value.contains(arg)
val sqlfunc = udf(func)
val filtered = xdf.filter(sqlfunc(col("name")))

xdf.show()
filtered.show()
Run Code Online (Sandbox Code Playgroud)

产量

name cnt
login1 192
login2 146
login3 72

名字c​​nt
login3 72


Iul*_*gos 12

  1. 你应该广播一个Set,而不是一个Array比线性更快的搜索.

  2. 您可以让Eclipse运行您的Spark应用程序.这是如何做:

正如邮件列表中所指出的,spark-sql假定其类由原始类加载器加载.如果Java和Scala库作为引导类路径的一部分加载,那么Eclipse中就不是这种情况,而用户代码及其依赖项则在另一个中.您可以在启动配置对话框中轻松修复它:

  • 从"Bootstrap"条目中删除Scala Library和Scala Compiler
  • 添加(作为外部jar)scala-reflect,scala-library并添加scala-compiler到用户条目.

对话框应如下所示:

在此输入图像描述

编辑:Spark错误已得到修复,不再需要此解决方法(因为v.1.4.0)