【发布时间】:2015-10-18 11:57:54
【问题描述】:
我正在使用火花流来处理文件流。多个文件成批到达并从所有文件中激发过程数据。 我的用途是获取进入后续批次的文件的每条记录的总和。例如:
- key: key_1 value: 10 --> batch1
- key: key_1 value: 05 --> batch1
- key: key_1 value: 19 --> batch2
- key: key_1 value: 11 --> batch3
- key: key_1 value: 10 --> batch4
我需要如下输出:
- 处理第一批后,我需要输出为 => key: key_1 val: 15
- 处理第二批后,我需要输出为 => key: key_1 val: 34
- 处理第三批后,我需要输出为 => key: key_1 val: 45
- 处理第 4 批后,我需要输出为 => key: key_1 val: 55
- 处理第 5 批后,我需要输出为 => key: key_1 val: 55
我的 reduceByKeyAndWindow() 代码如下:
JavaPairDStream<String, Summary> grpSumRDD = sumRDD.reduceByKeyAndWindow(GET_GRP_SUM, Durations.minutes(2*batchInterval), Durations.minutes(batchInterval));
private static final Function2<Summary, Summary, Summary> GET_GRP_SUM = new Function2<Summary, Summary, Summary>() {
private static final long serialVersionUID = 1L;
public Summary call(Summary s1, Summary s2) throws Exception {
try {
Summary s = new Summary();
long grpCnt = s1.getDelta() + s2.getDelta();
s.setDeltaSum(grpCnt);
return s;
} catch (Exception e) {
logger.error(" ==== error in CKT_GRP_SUM ==== :"+e);
return new Summary();
}
}
};
我从上面的实现中得到的输出如下:
- 处理第一批后,我得到输出 => key: key_1 value: 15
- 处理第二批后,我得到输出 => key: key_1 value: 34
- 处理第三批后,我得到输出 => key: key_1 value: 30
- 处理第 4 批后,我得到输出 => key: key_1 value: 21
- 处理第 5 批后,我得到输出 => key: key_1 value: 10
根据 reduceByKeyAndWindow() 的输出,它似乎正在计算先前批次数据和当前批次数据的聚合。 但我的要求是对上一批的聚合数据和当前的批数据进行聚合。这样根据上面的例子 它应该在第 4 和第 5 批结束时输出为 [(((15)+19)+11)+10 = 55]。
我读到了 reduceByKeyAndWindow() 和 invFunc 可以实现以获得预期的输出。我试图实现它类似于 GET_GRP_SUM 但它没有给我预期的结果。任何有关正确实施以获得所需输出的帮助将不胜感激。
我正在使用 java 1.8.45 和 spark 版本 1.4.1 和 hadoop 版本 2.7.1。
我用 reduceByKeyAndWindow() 对 invFunc 的实现
JavaPairDStream<String, Summary> grpSumRDD = sumRDD.reduceByKeyAndWindow(GET_GRP_SUM, INV_GET_GRP_SUM, Durations.minutes(2*batchInterval), Durations.minutes(batchInterval));
private static final Function2<Summary, Summary, Summary> INV_GET_GRP_SUM = new Function2<Summary, Summary, Summary>() {
private static final long serialVersionUID = 1L;
public Summary call(Summary s1, Summary s2) throws Exception {
try {
Summary s = new Summary();
long grpCnt = s1.getDelta() + s2.getDelta();
s.setDeltaSum(grpCnt);
return s;
} catch (Exception e) {
logger.error(" ==== error in INV_GET_GRP_SUM ==== :"+e);
return new Summary();
}
}
};
我已经像上面那样实现了我的 invFunc,这并没有给我预期的输出。我这里分析的是 s1 和 s2 给我之前批次的聚合值,我觉得我不太确定。
我尝试更改我的 invFunc 实现,如下所示:
private static final Function2<Summary, Summary, Summary> INV_GET_GRP_SUM = new Function2<Summary, Summary, Summary>() {
private static final long serialVersionUID = 1L;
public Summary call(Summary s1, Summary s2) throws Exception {
try {
return s1;
} catch (Exception e) {
logger.error(" ==== error in INV_GET_GRP_SUM ==== :"+e);
return new Summary();
}
}
};
这个实现给了我预期的输出。但我面临的问题是带有 invFunc 的 reduceByKeyAndWindow() 不会自动删除旧键。我又看了几篇文章,发现我需要编写自己的过滤器函数,该函数将删除具有 0 值(无值)的旧键。
再次,我不确定如何编写过滤函数来删除具有 0 值(无值)的旧键,因为我没有具体了解 s1 和 s2 返回到 INV_GET_GRP_SUM 的内容。
【问题讨论】:
标签: java apache-spark spark-streaming