【发布时间】:2016-04-06 14:39:37
【问题描述】:
我目前正在尝试用 Java 编写一个 Spark 作业,用于计算数据集中列的积分。
数据如下:
DateTime velocity (in km/h) vehicle
2016-03-28 11:00:45 80 A
2016-03-28 11:00:45 75 A
2016-03-28 11:00:46 70 A
2016-03-28 11:00:47 68 A
2016-03-28 11:00:48 72 A
2016-03-28 11:00:48 75 A
...
2016-03-28 11:00:47 68 B
2016-03-28 11:00:48 72 B
2016-03-28 11:00:48 75 B
要计算每条线路的距离(以公里为单位),我必须定义当前线路和下一条线路之间的时间差,并将其乘以速度。 然后必须将结果添加到上一行的结果中,以检索当时行驶的“总距离”。
我现在想出了这样的事情。但它会为每个地图作业计算一辆车,并且可能有数百万条记录......
final JavaRDD<String[]> input = sc.parallelize(Arrays.asList(
new String[]{"2016-03-28", "11:00", "80", "VIN1"},
new String[]{"2016-03-28", "11:00", "60", "VIN1"},
new String[]{"2016-03-28", "11:00", "50", "VIN1"},
new String[]{"2016-03-28", "11:01", "80", "VIN1"},
new String[]{"2016-03-28", "11:05", "80", "VIN1"},
new String[]{"2016-03-28", "11:09", "80", "VIN1"},
new String[]{"2016-03-28", "11:00", "80", "VIN2"},
new String[]{"2016-03-28", "11:01", "80", "VIN2"}
));
// grouping by vehicle and date:
final JavaPairRDD<String, Iterable<String[]>> byVinAndDate = input.groupBy(new Function<String[], String>() {
@Override
public String call(String[] record) throws Exception {
return record[0] + record[3]; // date, vin
}
});
// mapping each "value" (all record matching key) to result
final JavaRDD<String[]> result = byVinAndDate.mapValues(new Function<Iterable<String[]>, String[]>() {
@Override
public String[] call(Iterable<String[]> records) throws Exception {
final Iterator<String[]> iterator = records.iterator();
String[] previousRecord = iterator.next();
for (String[] record : records) {
// Calculate difference current <-> previous record
// Add result to new list
previousRecord = record;
}
return new String[]{
previousRecord[0],
previousRecord[1],
previousRecord[2],
previousRecord[3],
NewList.get(previousRecord[0]+previousRecord[1]+previousRecord[2]+previousRecord[2])
};
}
}).values();
我完全不知道如何将这个问题转化为映射/归约转换,同时又不失分布式计算的好处。
我知道这与 MR 和 Spark 的本质背道而驰,但任何有关如何互连数据行或以优雅方式解决此问题的建议都会非常有帮助:)
谢谢!
【问题讨论】:
标签: java hadoop apache-spark rdd integral