【发布时间】:2017-05-28 08:12:46
【问题描述】:
我有数据集,如果它满足某些状态,我需要计算数据的连续性。示例数据集如下。用例是,如果交换 id 连续处于风险和不稳定状态,则将该周的计数增加 1 并与数据集合并。我正在尝试使用 Spark。
Date Exchange Id Status Consecutiveness
5/05/2017 a RISKY 0
5/05/2017 b Stable 0
5/05/2017 c Stable 0
5/05/2017 d UNSTABLE 0
5/05/2017 e UNKNOWN 0
5/05/2017 f UNKNOWN 0
6/05/2017 a RISKY 1
6/05/2017 b Stable 0
6/05/2017 c Stable 0
6/05/2017 d UNSTABLE 1
6/05/2017 e UNSTABLE 1
6/05/2017 f UNKNOWN 0
我的方法如下。
- 为具有风险和不稳定的当前日期交换创建数据框 条件
- 为前一个日期创建另一个数据框,以便交易所具有 有风险且不稳定
- 加入 2 个数据框并获得不符合条件的交换
- 更新当前日期的连续性
- 与原始数据集合并。
我正在尝试执行以下命令。但是,遇到问题,无法继续进行 3,4,5
case class Telecom(Date: String, Exchange: String, Stability: String, Cosecutive: Int)
val emp1 = sc.textFile("file:/// Filename").map(_.split(",")).map(emp1=>Telecom(emp1(0),emp1(1),emp1(2),emp1(4).trim.toInt)).toDF()
val PreviousWeek = sqlContext.sql("select * from T1 limit 10")
emp1.registerTempTable("T1")
val FailPreviousWeek = sqlContext.sql("Select Exchange, Count from T1 where Date = '5/05/2017' and Stability in ('RISKY','UNSTABLE')")
val FailCurrentWeek = sqlContext.sql("Select Exchange, Count from T1 where Date = '6/05/2017' and Stability in ('RISKY','UNSTABLE')")
FailCurrentWeek.join(FailPreviousWeek, FailCurrentWeek("Exchange") === FailPreviousWeek("Exchange"))
val UpdateCurrentWeek = FailCurrentWeek.select($"Exchange",$"Count" +1)
Val UpdateDataSet = emp1.join(UpdateCurrentWeek)
val UpdateCurrentWeek = FailCurrentWeek.select($"Exchange".alias("Exchangeid"),$"Count" +1)
【问题讨论】:
-
为什么
e一周的6/05/2017是1?
标签: scala apache-spark apache-spark-sql