【问题标题】:Subtracting values from first and last row from Cassandra in Spark从 Spark 中的 Cassandra 中减去第一行和最后一行的值
【发布时间】:2016-05-23 16:00:47
【问题描述】:

我有这段代码,它从 Cassandra 获取 RDD,然后为每个键提取第一行和最后一行并减去它们。

val rdd = sc.cassandraTable("keyspace","table")
    .select("column1", "column2", "column3", "column4","column5")
    .as((i:String, p:String, e:String, c:Double, a:java.util.Date) => ((i), (c, a, p, e)))
    .groupByKey.mapValues(v => v.toList)
    .cache

val finalValues = rdd.mapValues(v => v.head)
val initialValues = rdd.mapValues(v => v.last)
val valuesCombined = finalValues.join(initialValues)

val results = valuesCombined.map(v => (v._2._1._1 - v._2._2._1))

它的性能好还是有更好的解决方案?我不确定是否将整个数据集缓存在内存中。

【问题讨论】:

  • 这不会先提取和最后提取。它只是提取任意行,恰好在 groupByKey 之后的第一个或最后一个?是你想要的吗?如果不是,你想如何排序这些值?
  • Cassandra 在插入期间按日期对表行进行排序。
  • groupByKey 不保证在随机播放期间会保留顺序。
  • 感谢您的信息。有什么快速的解决方案吗?

标签: scala apache-spark cassandra datastax-enterprise


【解决方案1】:

groupByKey 打乱数据,不再保证分组值的顺序。它也相当昂贵。

如果您真的想在RDDs 而不是DataFrames 上进行操作,并且订购是基于您可以使用aggregateByKey 的日期:

import scala.math.Ordering

type Record = (String, String, String, Double, java.util.Date)
val RecordOrd = Ordering.by[Record, java.util.Date](_._5)

val minRecord = ("", "", "", 0.0, new java.util.Date(Long.MinValue))
val maxRecord = ("", "", "", 0.0, new java.util.Date(Long.MaxValue))

def minMax(x: (Record, Record), y: (Record, Record)) = {
  (RecordOrd.min(x._1, y._1), RecordOrd.max(x._2, y._2))
}

rdd.aggregateByKey((maxRecord, minRecord))(
  (acc, x) => minMax(acc, (x, x)),
  minMax
)

使用DataFrames,您可以尝试以下操作:

import org.apache.spark.sql.functions.{col, lag, lead, when, row_number, max}
import org.apache.spark.sql.expressions.Window

val partition = Seq("column1")
val order = Seq("column5")
val columns = Seq("column2", "column3", "column4","column5")

val w = Window
  .partitionBy(partition.head, partition.tail: _*)
  .orderBy(order.head, order.tail: _*)

// Lead / lag of row number to mark first / last row in the group
val rn_lag = lag(row_number.over(w), 1).over(w)
val rn_lead = lead(row_number.over(w), 1).over(w)

// Select value if first / last row in the group otherwise null
val firstColumns = columns.map(
  c => when(rn_lag.isNull, col(c)).alias(s"${c}_first"))
val lastColumns = columns.map(
  c => when(rn_lead.isNull, col(c)).alias(s"${c}_last"))

// Add columns with first / last vals
val expanded = df.select(
  partition.map(col(_)) ++ firstColumns ++ lastColumns: _*)

// Aggregate to drop nulls
val aggExprs = expanded.columns.diff(partition).map(c => max(c).alias(c))
expanded.groupBy(partition.map(col(_)): _*).agg(aggExprs.head, aggExprs.tail: _*)

还有一些其他方法可以使用DataFrames 解决此问题,包括通过structsDataSet API 进行排序。查看我对SPARK DataFrame: select the first row of each group的回复

【讨论】:

  • 感谢您的意见。我已经在 Spark 上完成了 Datastax 教程,它基于 RDD,最后只提到了 DataFrame,我认为 RDD 是要走的路。现在在阅读了link 之后,我知道 DataFrame 在性能方面更好。我将尝试使用 DataFrames 编写代码。如果您有时间,我将不胜感激使用 DataFrames 编写它。
  • @PawełSzychiewicz 我已经将另一个答案与一些示例联系起来,您可以如何解决组中的第一(最后)行选择,这可能比窗口函数更直观。
【解决方案2】:

首先 - 我假设 all 变量是指名为 rdd 的变量?创建后,您不需要使用 join(这在性能方面代价高昂),您可以简单地将每个元素直接映射到您需要的结果:

val results = all.mapValues(v => v.head - v.last).values

现在 - 因为我们只对 RDD 执行了一个操作,所以我们也可以摆脱 cache()

【讨论】:

  • 这行不通。 headlast 可以是任意元素,尤其是因为 groupByKey 禁用了地图端聚合。如果你想mapValues你应该先执行命令。
猜你喜欢
  • 2019-07-24
  • 2018-03-07
  • 2022-01-15
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多