【问题标题】:Flink: enrich a data set with a new column based on some computationFlink:基于一些计算用新列丰富数据集
【发布时间】: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


    【解决方案1】:

    这是一个多线程问题。你有多少个任务槽?

    我必须在运行之前清理您的代码 - 我建议您在以后发布完整的工作示例,以便您有机会获得更多答案。

    您跟踪计数的方式不是线程安全的,因此如果您有多个任务槽,您将遇到计数值不准确的问题。

    如数据工匠字数统计示例所示,正确的计数方法是使用元组中的第三个槽来简单地存储值 1,然后对数据集求和。

    resultTuple = new Tuple3(m.f0,m.f1, 1L);
    

    然后

    csvOut2.sum(2).print();
    

    其中 2 是包含值 1 的元组的索引。

    【讨论】:

    • 谢谢!我正在使用带有 1 个任务槽的 flink 的本地/独立部署(我在作业管理器 GUI 中看到)。我认为这可能是线程问题,所以我也尝试了AtomicLong。此外,计算元组很简单。我目前可以做到。我的问题是在每行的末尾添加一列,其中包含当前计数。
    • 对不起,我之前看错了。可能您需要使用状态,因为值取决于之前处理的内容。或者你可以使用累加器dataartisans.github.io/flink-training/dataSet/2-slides.html
    • 其实累加器在这里不合适。
    • 确实,在处理每条记录时,我需要一些方法来保持状态,因为当前记录依赖于它之前的每条记录。
    • 状态仅在键控流上可用
    猜你喜欢
    • 1970-01-01
    • 2011-07-27
    • 1970-01-01
    • 2011-08-19
    • 1970-01-01
    • 1970-01-01
    • 2011-02-11
    • 1970-01-01
    相关资源
    最近更新 更多