【问题标题】:How to sort data in spark streaming如何在火花流中对数据进行排序
【发布时间】:2015-01-06 09:04:52
【问题描述】:

我是 spark 新手,并尝试编写一些基于 spark 和 spark 流的示例代码。

到目前为止,我已经在spark中实现了排序功能,代码如下:

  def sort(listSize: Int, slice: Int): Unit = {
    val conf = new SparkConf().setAppName(getClass.getName)
    val spark = new SparkContext(conf)
    val data = genRandom(listSize)
    val distData = spark.parallelize(data, slice)
    val result = distData.sortBy(x => x, true)
    val finalResult = result.collect()
    val step5 = System.currentTimeMillis()
    printlnArray(finalResult, 0, 10)
    spark.stop()
  }

  /**
   * generate random number
   * @return
   */
  def genRandom(listSize: Int): List[Int] = {
    val range = 100000
    var listBuffer = new ListBuffer[Int]
    val random = new Random()
    for (i <- 1 to listSize) listBuffer += random.nextInt(range)
    listBuffer.toList
  }

  def printlnArray(list: Array[Int], start: Int, offset: Int) {
    for (i <- start until start + offset) println(">>>>>>>>> list : " + i + " | " + list(i))
  }

我在火花流上实现排序功能时遇到了麻烦。据我所知,spark RDD在spark core中提供了sort API,但是spark streaming中没有这样的API,有人知道怎么做吗?谢谢

这是一个转储问题,但在网上谷歌后,我没有找到正确的答案。如果有人知道如何解决,谢谢您的帮助。

【问题讨论】:

  • 您要对流中的每个microbatch 进行排序还是要对整个流进行排序?后者 - 就一般的流处理而言 - afaik 是不可能的。

标签: scala apache-spark


【解决方案1】:

您可以利用 DStream 的转换功能,通过使用底层 RDD 对其进行转换。

例如

myDStream.transform(rdd => rdd.sortByKey())

【讨论】:

  • 谢谢@Hawk66!想知道我们如何在 DStreams 上执行单一操作?比如说,如果我们只想要每个微批次中最顶层的条目。 myDStream.transform(rdd=&gt;rdd.sortByKey()).top(1) 或 myDStream.transform(rdd=&gt;rdd.sortByKey()).transform(rdd=&gt;rdd.top(1))?到目前为止没有任何效果
  • 没关系。在这个 SO 链接stackoverflow.com/questions/41483746/… 中得到了我的答案
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2021-06-13
  • 2012-12-06
  • 2020-10-05
  • 2018-02-11
  • 2017-01-29
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多