【发布时间】:2017-01-19 17:52:30
【问题描述】:
我正在尝试对数据集进行简单处理。
考虑一个包含两列String 类型的数据集。对于这个数据集,我想添加第三列Long,它累积了迄今为止在数据集中看到的记录数。
例子:
输入:
a,b
b,c
c,d
输出:
a,b,1
b,c,2
c,d,3
我尝试了以下解决方案,但得到了一个奇怪的结果:
DataSet<Tuple2<String, String>> csvInput = env.readCsvFile("src/main/resources/data_file")
.ignoreFirstLine()
.includeFields("11")
.types(String.class,String.class);
long cnt=0;
DataSet<Tuple3<String, String, Long>> csvOut2 = csvInput.map(new MyMapFunction(cnt));
private static class MyMapFunction implements MapFunction<Tuple2<String, String>, Tuple3<String, String, Long>> {
long cnt;
public MyMappingFunction(long cnt) {
this.cnt = cnt;
}
@Override
public Tuple3<String, String, Long> map(Tuple2<String, String> m) throws Exception {
Tuple3 <String ,String, Long> resultTuple = new Tuple3(m.f0,m.f1, Long.valueOf(cnt));
cnt++;
return resultTuple;
}
}
当我将此解决方案应用于具有 100 个条目的文件时,我得到的计数是 47 而不是 100。计数器在 53 处重新启动。同样,当我将它应用于更大的文件时,计数器会不时重置为时间所以我没有得到总行数。
您能否解释一下为什么我的实现会以这种方式运行?另外,我的问题有什么可能的解决方案?
谢谢!
【问题讨论】:
标签: count dataset apache-flink