【问题标题】:Lookup in Spark dataframes在 Spark 数据帧中查找
【发布时间】:2017-05-07 15:07:37
【问题描述】:

我使用的是 Spark 1.6,我想知道如何在数据帧中实现查找。

我有两个数据框员工和部门。

  • 员工数据框

    -------------------
    Emp Id | Emp Name
    ------------------
    1 | john
    2 | David
    
  • 部门数据框

    --------------------
    Dept Id | Dept Name | Emp Id
    -----------------------------
    1 | Admin | 1
    2 | HR | 2
    

我想从员工表中查找 emp id 到部门表并获取部门名称。因此,结果集将是

Emp Id | Dept Name
-------------------
1 | Admin
2 | HR

如何在 SPARK 中实现这个查找 UDF 功能。我不想在两个数据帧上都使用 JOIN。

【问题讨论】:

  • 你需要的是“加入”两个数据帧......如果一个非常小,请使用广播加入。
  • 到目前为止你有什么代码示例吗?
  • 这就是join 的意思。尝试通过其他方式实现这一点的原因是什么?
  • 我已经使用join实现了。但我也想探索使用查找概念(只是为了学习并了解实现和性能的差异)。有人知道吗?
  • 您可以将您的数据帧转换为PairRDD(例如,提供lookup 方法的RDD[(Int,String)])

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


【解决方案1】:

正如 cmets 中已经提到的,加入数据框是可行的方法。

您可以使用查找,但我认为没有“分布式”解决方案,即您必须将查找表收集到驱动程序内存中。另请注意,此方法假定 EmpID 是唯一的:

import org.apache.spark.sql.functions._
import sqlContext.implicits._
import scala.collection.Map

val emp = Seq((1,"John"),(2,"David"))
val deps = Seq((1,"Admin",1),(2,"HR",2))

val empRdd = sc.parallelize(emp)
val depsDF = sc.parallelize(deps).toDF("DepID","Name","EmpID")


val lookupMap = empRdd.collectAsMap()
def lookup(lookupMap:Map[Int,String]) = udf((empID:Int) => lookupMap.get(empID))

val combinedDF = depsDF
  .withColumn("empNames",lookup(lookupMap)($"EmpID"))

我最初的想法是将empRdd 传递给UDF 并使用PairRDD 上定义的lookup 方法,但这当然行不通,因为您不能在转换中使用火花动作(即lookup) (即 UDF)。

编辑:

如果您的 empDf 有多个列(例如 Name、Age),您可以使用此

val empRdd = empDf.rdd.map{row =>
      (row.getInt(0),(row.getString(1),row.getInt(2)))}


    val lookupMap = empRdd.collectAsMap()
    def lookup(lookupMap:Map[Int,(String,Int)]) =
         udf((empID:Int) => lookupMap.lift(empID))

    depsDF
      .withColumn("lookup",lookup(lookupMap)($"EmpID"))
      .withColumn("empName",$"lookup._1")
      .withColumn("empAge",$"lookup._2")
      .drop($"lookup")
      .show()

【讨论】:

  • 谢谢。在您的示例中, empRdd 不是数据框。如何将我的数据框转换为 rdd 以使用 collectAsMap 函数?
  • @Prasan 在数据帧上有一个rdd 方法,所以你必须像val empRdd = empDf.rdd.map(row => (row.getInt(0),row.getString(1)))这样的东西
  • 是的,但是 collectAsMap 不是 RDD 的成员。当我尝试使用 collect 而不是 collectAsMap 时,它给了我一个无法在查找函数中使用的数组。
  • @Prasan 如果您的 RDD 是 RDD[(Int,String)] 类型,那么您应该能够通过隐式转换使用 collectAsMapPairRDDFunctions
  • 有没有办法返回多列而不是只返回一列?假设您的 emp 数据框示例中也有 emp_age,是否可以使用此查找概念来显示它?
【解决方案2】:

正如您所说,您已经拥有 Dataframe,那么按照以下步骤操作非常容易:

1)创建一个sqlcontext

val sqlContext = new org.apache.spark.sql.SQLContext(sc)

2) 为所有 3 个创建临时表 例如:

EmployeeDataframe.createOrReplaceTempView("EmpTable")

3) 使用 MySQL 查询进行查询

val MatchingDetails = sqlContext.sql("SELECT DISTINCT E.EmpID, DeptName FROM EmpTable E inner join DeptTable G on " +
  "E.EmpID=g.EmpID")

【讨论】:

  • 感谢您的回复。正如我在帖子中提到的,我不想使用 JOIN。
  • 为什么你不能使用 join 来达到目的
  • 这只是我发布的一个例子。在我的实时场景中,我们正在尝试将复杂的 Informatica 映射转换为 SPARK。我正在尝试在 SPARK 中复制查找功能(就像 informatica 一样)。
  • @Prasan Informatica 是一个旧产品。从它迁移时,不要试图保留/复制“旧方式”。 Spark 中的lookup 非常慢,而 join 在其不同的风格(内部、外部、左外部......)中非常优化。对于您使用的示例,join 是 Spark 的方式。如果您有其他用例,请在另一个问题上提供足够的上下文。
  • 我已经使用join实现了。但我也想探索使用查找概念(只是为了学习并查看实现和性能的差异)
【解决方案3】:

从一些“查找”数据开始,有两种方法:

方法 #1 -- 使用查找数据帧

// use a DataFrame (via a join)
val lookupDF = sc.parallelize(Seq(
  ("banana",   "yellow"),
  ("apple",    "red"),
  ("grape",    "purple"),
  ("blueberry","blue")
)).toDF("SomeKeys","SomeValues")

方法 #2 -- 在 UDF 中使用地图

// turn the above DataFrame into a map which a UDF uses
val Keys = lookupDF.select("SomeKeys").collect().map(_(0).toString).toList
val Values = lookupDF.select("SomeValues").collect().map(_(0).toString).toList
val KeyValueMap = Keys.zip(Values).toMap

def ThingToColor(key: String): String = {
  if (key == null) return ""
  val firstword = key.split(" ")(0) // fragile!
  val result: String = KeyValueMap.getOrElse(firstword,"not found!")
  return (result)
}

val ThingToColorUDF = udf( ThingToColor(_: String): String )

获取将要查找的事物的样本数据框:

val thingsDF = sc.parallelize(Seq(
  ("blueberry muffin"),
  ("grape nuts"),
  ("apple pie"),
  ("rutabaga pudding")
)).toDF("SomeThings")

方法#1是加入lookup DataFrame

这里,rlike 正在进行匹配。 null 出现在不起作用的地方。查找 DataFrame 的两列都被添加。

val result_1_DF = thingsDF.join(lookupDF, expr("SomeThings rlike SomeKeys"), 
                     "left_outer")

方法#2是使用UDF添加一列

这里只添加了 1 列。并且 UDF 可以返回非 Null 值。但是,如果查找数据非常大,它可能无法按要求“序列化”以发送给集群中的工作人员。

val result_2_DF = thingsDF.withColumn("AddValues",ThingToColorUDF($"SomeThings"))

这给了你:

在我的例子中,我有一些超过 100 万个值的查找数据,所以方法 #1 是我唯一的选择。

【讨论】:

  • 感谢warrens 这个“val KeyValueMap = Keys.zip(Values).toMap”如何保证只有“正确”的值被映射到“正确”的键?
  • 有关如何在 UDF 中的地图中处理此问题的任何线索 stackoverflow.com/questions/63935600/…
猜你喜欢
  • 2018-06-07
  • 1970-01-01
  • 1970-01-01
  • 2021-07-17
  • 1970-01-01
  • 2022-01-08
  • 1970-01-01
  • 1970-01-01
  • 2016-01-03
相关资源
最近更新 更多