【问题标题】:Custom scalding tap (or Spark equivalent)自定义烫伤水龙头(或 Spark 等效)
【发布时间】:2014-06-07 05:10:14
【问题描述】:

我正在尝试使用自定义文件格式转储 Hadoop 集群(通常在 HBase 中)上的一些数据。

我想做的或多或少如下:

  • 从分布式记录列表开始,例如 Scalding pipe 或类似物
  • 通过某些计算函数对项目进行分组
  • 使属于同一组的项目驻留在同一服务器上
  • 在每个组上,应用转换 - 包括排序 - 并将结果写入磁盘。其实我需要写一堆MapFile——本质上是排序的SequenceFile,加上一个索引。

我想用 Scalding 来实现上面的,但我不知道如何做最后一步。

当然,虽然不能以分布式方式写入已排序的数据,但将数据拆分为块然后写入本地排序的每个块应该仍然可行。尽管如此,我还是找不到任何用于 map-reduce 作业的 MapFile 输出实现。

我认识到对非常大的数据进行排序是个坏主意,这就是我计划在单个服务器上将数据拆分成块的原因。

有没有办法用 Scalding 做类似的事情?可能我可以直接使用 Cascading 或其他管道框架,例如 Spark。

【问题讨论】:

    标签: scala hadoop cascading apache-spark scalding


    【解决方案1】:

    使用 Scalding(和底层 Map/Reduce),您将需要使用 TotalOrderPartitioner,它会进行预采样以创建输入数据的适当存储桶/拆分。

    由于磁盘数据的访问路径更快,使用 Spark 会加快速度。但是,它仍然需要对磁盘/hdfs 进行洗牌,因此不会好几个数量级。

    在 Spark 中,您将使用 RangePartitioner,它获取分区数和 RDD:

    val allData = sc.hadoopRdd(paths)
    val partitionedRdd = sc.partitionBy(new RangePartitioner(numPartitions, allData)
    val groupedRdd = partitionedRdd.groupByKey(..).
    // apply further transforms..
    

    【讨论】:

      猜你喜欢
      • 2014-06-19
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2013-06-12
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多