【问题标题】:How to Sort DataFrame with my Comparator using Scala?如何使用 Scala 使用我的 Comparator 对 DataFrame 进行排序?
【发布时间】:2019-08-03 07:44:54
【问题描述】:

我想使用我自己的比较器根据列对 DataFrame 进行排序。在 Spark SQL 中可以做到这一点吗?

例如,假设我有一个注册为表“MyTable”的DataFrame,其中有一列“Day”,其类型为“string”:

id  | Day  
--------------------
1   | Fri           
2   | Mon           
3   | Sat           
4   | Sun           
5   | Thu           

我想执行这个查询:

SELECT * FROM MyTable ORDER BY Day

我想用我自己的比较器订购“日”列。我考虑过使用UDF,但我不知道是否可能。请注意,我真的想在 Sort/Order By 操作中使用我的比较器。我不想将字符串从 Day 列转换为 Datetime 或类似的东西。

【问题讨论】:

  • 在 spark 中无法使用通用的 java 比较器进行排序。您需要定义一个排序键(该类型是按 s.a. long、string、date... 排序的类型)并为数据集使用 orderBy 或为 RDD 使用 sortBy。如果您告诉我们您的具体逻辑,也许我们可以考虑一个适合您的解决方案。
  • 感谢您的回答。我必须做两件事:1)使用本例中的 SQL 查询; 2) 当我必须执行 ORDER BY 时,使用运算符对某个列进行排序。一旦无法在 UDF 中使用通用比较器,您会推荐我另一种选择吗?我考虑过使用 Spark 规则/策略将 DF 转换为 RDD 并将 SortBy 与我的比较器一起使用。但我不知道具体该怎么做。
  • 其实我有点错。可以将比较器(scala 中的排序)与 RDD 一起使用。我添加了一个解决方案来解释如何做到这一点。我也谈论替代品。

标签: scala sorting apache-spark apache-spark-sql


【解决方案1】:

这是使用数据框的一般方法

val df = spark.sql("SELECT * FROM MyTable")

df.orderby("yourcolumn")

orderby docs


如果您的数据较少(似乎您只有周名),那么您可以收集为列表并使用 scala sortWith 函数

sortWith 函数根据比较对该序列进行排序 功能。它需要一个比较器功能并根据它进行排序。你 可以提供自己的自定义比较功能。

和你的例子不同:

scala> case class Emp(id: Int, name: String, salary: Double)
defined class Emp

scala> val emp1 = Emp(1, "james", 13000.00)
emp1: Emp = Emp(1,james,13000.0)

scala> val emp2 = Emp(2, "michael", 12000.00)
emp2: Emp = Emp(2,michael,12000.0)

scala> val emp3 = Emp(3, "Ram", 15000.00)
emp3: Emp = Emp(3,Ram,15000.0)

scala> val empList = List(emp1,emp2,emp3)
empList: List[Emp] = List(Emp(1,james,13000.0), Emp(2,michael,12000.0), Emp(3,Ram,15000.0))

// sort in descending order on the basis of salary.
scala> empList.sortWith(_.salary > _.salary)

其他选项是: How to sort an RDD in Scala Spark? 要使用此选项,您需要将数据框转换为 PairedRDD,然后使用那里给出的答案进行排序。

【讨论】:

  • 感谢您的回答!它帮助我更好地理解这些概念。正如我上面所说,我需要像我展示的示例中那样执行 SQL 查询。您知道是否有任何方法可以使用 Spark Rules/Spark Strategies 来转换 Spark 计划,将 DataFrame 转换为 DataSet/RDD 并使用比较器?在我看来,这是唯一可能的解决方案。
  • 用普通的 sql 是不可能的
【解决方案2】:

在 SparkSQL 中,您别无选择,需要将orderBy 与一个或多个列一起使用。使用 RDD,如果您愿意,可以使用自定义的类似 java 的比较器。事实上,这里是 RDD (cf the scaladoc of Spark 2.4) 的 sortBy 方法的签名:

def sortBy[K](f: (T) ⇒ K, ascending: Boolean = true, numPartitions: Int = this.partitions.length)
    (implicit ord: Ordering[K], ctag: ClassTag[K]): RDD[T] 

这意味着您可以提供您选择的Ordering,这与java Comparator 完全相同(Ordering 实际上继承自Comparator)。

为简单起见,假设我想按列“x”的绝对值排序(这可以在没有比较器的情况下完成,但假设我需要使用比较器)。我首先在行上定义比较器:

class RowOrdering extends Ordering[Row] {
    def compare(x : Row, y : Row): Int = x.getAs[Int]("x").abs - y.getAs[Int]("x").abs
}

现在让我们定义数据并对其进行排序:

val df = Seq( (0, 1),(1, 2),(2, 4),(3, 7),(4, 1),(5, -1),(6, -2),
    (7, 5),(8, 5), (9, 0), (10, -9)).toDF("id", "x")
val rdd = df.rdd.sortBy(identity)(new RowOrdering(), scala.reflect.classTag[Row])
val sorted_df = spark.createDataFrame(rdd, df.schema)
sorted_df.show
+---+---+
| id|  x|
+---+---+
|  9|  0|
|  0|  1|
|  4|  1|
|  5| -1|
|  6| -2|
|  1|  2|
|  2|  4|
|  7|  5|
|  8|  5|
|  3|  7|
| 10| -9|
+---+---+

另一种解决方案是定义隐式排序,以便在排序时无需提供它。

implicit val ord = new RowOrdering()
df.rdd.sortBy(identity)

最后,请注意df.rdd.sortBy(_.getAs[Int]("x").abs) 将获得相同的结果。此外,您可以使用元组排序来做更复杂的事情,例如按绝对值排序,如果相等,则将正值放在首位:

df.rdd.sortBy(x => (x.getAs[Int]("x").abs, - x.getAs[Int]("x"))) //RDD
df.orderBy(abs($"x"), - $"x") //dataframe

【讨论】:

  • 如果我也想使用比较器进行 GroupBy 操作,过程基本相同,对吧?将 DataFrame 转换为 RDD 并使用: groupBy[K](f: (T) ⇒ K, p: Partitioner)(implicit kt: ClassTag[K], ord: Ordering[K] = null): RDD[(K, Iterable [T])]
  • 从未尝试过,但考虑到签名似乎是可能的。不过,您将需要定义一个兼容的分区器,即始终将相等的键放在同一个分区中的分区器。
  • 我试过了,你可以让它工作,但它不会对不相等的键进行分组,即使你的命令说它们是一致的分区器。不过,您的密钥将被正确排序。根据您的用例,这可能就足够了。
  • 我为 Group By 示例发布了一个新的question。它可以帮助其他开发人员。感谢您的帮助!
猜你喜欢
  • 1970-01-01
  • 2019-04-07
  • 1970-01-01
  • 2021-11-14
  • 1970-01-01
  • 2011-02-19
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多