【问题标题】:How to output data with Hive-style directory structure in Scalding?如何在 Scalding 中以 Hive 样式的目录结构输出数据?
【发布时间】:2015-04-25 14:07:37
【问题描述】:

我们使用 Scalding 进行 ETL 并将输出生成为带有分区的 Hive 表。因此,我们希望分区的目录名称类似于“state=CA”。我们使用 TemplatedTsv 如下:

pipe
   // some other ETL
   .map('STATE -> 'hdfs_state) { state: Int => "State=" + state }
   .groupBy('hdfs_state) { _.pass }
   .write(TemplatedTsv(baseOutputPath, "%s", 'hdfs_state,
          writeHeader = false,
          sinkMode = SinkMode.UPDATE,
          fields = ('all except 'hdfs_state)))

我们采用来自How to bucket outputs in Scalding 的代码示例。 以下是我们遇到的两个问题:

  • except IntelliJ 无法解决:我错过了一些导入吗?我们不想在“fields = ()”语句中显式输入所有字段,因为字段是从 groupBy 语句中的代码派生的。如果明确输入,它们很容易不同步。
  • 这种方法看起来太老套了,因为我们正在创建一个额外的列,以便 Hive/Hcatalog 可以处理目录名称。我们想知道实现它的正确方法应该是什么?

非常感谢!

【问题讨论】:

    标签: scalding


    【解决方案1】:

    抱歉,前面的示例是伪代码。下面我将给出一个带有输入数据示例的小代码。

    请注意,这仅适用于 Scalding 0.12.0 或更高版本

    让我们想象一下我们输入的如下定义一些购买数据的图像,

    user1   1384034400  6   75
    user1   1384038000  6   175
    user2   1383984000  48  3
    user3   1383958800  48  281
    user3   1384027200  9   7
    user3   1384027200  9   11
    user4   1383955200  37  705
    user4   1383955200  37  15
    user4   1383969600  36  41
    user4   1383969600  36  21
    

    制表符分隔,第 3 列是州编号。这里我们有整数,但对于基于字符串的状态,您可以轻松适应。

    此代码将读取输入并将它们放入“State=stateid”输出文件夹存储桶中。

    class TemplatedTsvExample(args: Args) extends Job(args) {
    
      val purchasesPath = args("purchases")
      val outputPath    = args("output")
    
      // defines both input & output schema, you can also make separate for each of them
      val ioSchema = ('USERID, 'TIMESTAMP, 'STATE, 'PURCHASE)
    
      val Purchases =
         Tsv(purchasesPath, ioSchema)
         .read
         .map('STATE -> 'STATENAME) { state: Int => "State=" + state } // here you can make necessary changes
         .groupBy('STATENAME) { _.pass } // this is optional
         .write(TemplatedTsv(outputPath, "%s", 'STATENAME, false, SinkMode.REPLACE, ioSchema))
    } 
    

    我希望这会有所帮助。有什么不清楚的可以问我。

    你可以找到完整代码here.

    【讨论】:

    • 谢谢莫拉佐!我需要遵循其分区命名约定的 hive 目录结构,而不明确列出所有字段。如果我这样做,` .groupBy('STATENAME) {_.size('user_count)..max('PURCHASE -> 'max_purchase)} ` 请注意,我们需要文件中的 'STATE 值,而 'STATENAME 是目录。我们可以明确列出,如您的示例所示。但我正在寻找代码,只是排除了仅为目录命名目的而创建的字段,即本例中的“STATENAME”。
    • 在这种情况下,只需将 'STATE 添加到您的分组中,` .groupBy('STATENAME, 'STATE){_.size('user_count)..max('PURCHASE -> 'max_purchase)} .write(TemplatedTsv(outputPath, "%s", 'STATENAME, false, SinkMode.REPLACE, ('STATE, 'user_count, 'max_purchase)))` 分组只会保持分组和聚合字段,所以你可以添加'STATE分组。因为 'STATE 和 'STATENAME 是一对一的映射,所以不会改变逻辑。
    • 谢谢莫拉佐!有没有办法缩小你的伪代码: fields = ('all except 'STATENAME) 而不明确列出字段?
    • 嘿,欢迎您!我不知道。就算有,我也不知道。
    猜你喜欢
    • 2014-12-06
    • 1970-01-01
    • 2022-01-21
    • 1970-01-01
    • 2016-02-14
    • 1970-01-01
    • 1970-01-01
    • 2010-09-20
    • 1970-01-01
    相关资源
    最近更新 更多