【问题标题】:Bloating a dataset with spark & scala用 spark 和 scala 膨胀数据集
【发布时间】:2016-11-04 22:40:23
【问题描述】:

这是我的要求

输入

customer_id status  start_date  end_date
1   Y   20140101    20140105
2   Y   20140201    20140203

输出

customer_id status  date
1   Y   20140101
1   Y   20140102
1   Y   20140103
1   Y   20140104
1   Y   20140105
2   Y   20140201
2   Y   20140202
2   Y   20140202

我正在尝试使用 spark 中的笛卡尔积来实现这一点,但它看起来效率非常低。我的数据集太大了。我正在寻找更好的选择。

【问题讨论】:

  • 所以您希望每天在开始和结束之间有一条线,包括每个客户 ID。状态总是 Y 吗?如果不是,非 Y 的输出应该是什么?
  • status 也可以是 'N'.. 它是 Y 或 N

标签: scala apache-spark


【解决方案1】:

如果我的想法正确,你可以这样做:

  val conf = new SparkConf().setMaster("local[2]").setAppName("test")
  val sc = new SparkContext(conf)

  case class Input(customerId: Long, status: String, startDate: LocalDate, endDate: LocalDate)
  case class Output(customerId: Long, status: String, date: LocalDate)

  val input: RDD[Input] = sc.parallelize(Seq(
    Input(1, "Y", LocalDate.of(2014, 1, 1), LocalDate.of(2014, 1, 5)),
    Input(2, "Y", LocalDate.of(2014, 1, 1), LocalDate.of(2014, 1, 3))
  ))

  val result: RDD[Output] = input flatMap { input =>
    import input._
    val dates = Stream.iterate(startDate)(_.plusDays(1)).takeWhile(!_.isAfter(endDate))
    dates.map(date => Output(customerId, status, date))
  }

  result.collect().foreach(println)

输出:

Output(1,Y,2014-01-01)
Output(1,Y,2014-01-02)
Output(1,Y,2014-01-03)
Output(1,Y,2014-01-04)
Output(1,Y,2014-01-05)
Output(2,Y,2014-01-01)
Output(2,Y,2014-01-02)
Output(2,Y,2014-01-03)

【讨论】:

  • dates 将是一个普通的集合,但不是 RDD。所以可能不适合任何一个节点。
  • 我不认为日期集合可以大于 1 兆字节,所以这应该不是问题
  • 是的,抱歉,我错过了输入 flatMap,这应该意味着你又得到了一个 RDD。
  • 不幸的是,它在读取输入数据集本身时卡住了......在我的情况下,输入是 S3 中的文本文件......
猜你喜欢
  • 2018-02-10
  • 2018-09-26
  • 2013-07-15
  • 2014-07-25
  • 1970-01-01
  • 2021-10-12
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多