【问题标题】:Apache beam/Google dataflow enrich upstream records with downstream aggregatesApache Beam/Google 数据流使用下游聚合丰富上游记录
【发布时间】:2019-07-18 14:00:22
【问题描述】:

我创建了一个 Java apache 光束流管道,我计划在谷歌数据流上运行。它接收类似于以下内容的元素:

ipAddress, serviceUsed, errorOrSuccess, time, parameter, etc.

例如

'237.98.58.248', 'service1', 'error', '12345', 'randomParameter', etc

我目前根据事件时间将此数据窗口化到固定窗口中。我想使用我的管道计算每个窗口接收到的每个 ip 地址的错误和成功次数,然后丰富原始数据。

我希望调整每个原始元素以输出类似于以下内容的最终元素:

totalErrorsInThisWindow, totalSuccessInThisWindow, ipAddress, serviceUsed, errorOrSuccess, time, parameter, etc.

例如

'237.98.58.248', 'service1', 'error', '12345', 'randomParameter', etc
'149.142.114.250', 'service2', 'success', '12346', 'randomParameter', etc
'237.98.58.248', 'service3', 'error', '12344', 'randomParameter', etc
...

变成类似

'100', '1000', '237.98.58.248', 'service1', 'error', '12345', 'randomParameter', etc
'11', '34', '149.142.114.250', 'service2', 'success', '12346', 'randomParameter', etc
'100', '1000', '237.98.58.248', 'service3', 'error', '12344', 'randomParameter', etc
...

关于如何做到这一点的任何建议?

我知道几种方法如何在每个客户端、每个窗口的基础上计算 totalErrorsInThisWindowtotalSuccessInThisWindow - 一种方法是删除除 ipAddresserrorOrSuccess 之外的所有列,然后执行apply(Count.<String>perElement());。但是,我正在努力丰富原始数据。第一个想法是使用侧面输入,但我认为使用不断变化的侧面输入效果不佳。

另一个选项是为成功和失败维护一个基于键的状态变量,我可以在处理每个元素时增加它,并使用它来丰富同一个 DoFn 中的数据。但是,我遇到的问题是,只有在窗口中为每个键处理的最后一个元素才会具有正确的成功和失败值。

这是我可以用 state 做什么与我想用 state 做什么的一个例子:

输入:

'a'
'b'
'a'
'a'

使用状态可以得到的输出:

'a':1
'b':1
'a':2
'a':3

我想使用状态得到的输出:

'a':3
'b':1
'a':3
'a':3

我希望我的问题很清楚,我希望我目前的方法和挑战也很清楚。任何建议将不胜感激。

【问题讨论】:

    标签: google-cloud-dataflow apache-beam


    【解决方案1】:

    请查看 GroupByKey 和 Combiners,以及 how to use it with windowing

    我认为这样的事情会很好用。您可以按 IP 分组,应用窗口并计算错误和成功。

    PCollection<MyRecord> records = <Read from your source>
    
    PCollection<KV<string, MyRecord>> withIP = records.apply(ParDo.of(
        new DoFn<MyRecord, KV<KV<string, MyRecord>>>() {
          // Implement processElement and call outputWithTimestamp
        }
    ));
    
    
    PCollection<MyRecord> windowed = withIP.apply(
        Window.<MyRecord>into(FixedWindows.of(Duration.standardSeconds(60))));
    
    PCollection<KV<String, Iterable<MyRecords>>> grouped = 
        windowed.apply(GroupByKey.<String, MyRecords>create());
    
    
    PCollection<KV<String, MyErrorStats>> errorsPerIP =
      playerAccuracy.apply(Combine.<String, MyRecord, MyErrorStats>perKey(
        new MyErrorStatsCombiner())));
    
    
    public static class MyErrorStatsCombiner implements     
    SerializableFunction<Iterable<MyRecords>, MyErrorStats> {
      @Override
      public Integer apply(Iterable<MyRecords> record) {
        MyErrorStats stats = new MyErrorStats();
        for (int item : record) {
          stats.errorsInThisWindow += item.errorsInThisWindow;
          stats.successInThisWindow += item.successInThisWindow;
        }
        return stats;
      }
    }
    

    至于保留记录中的其他元数据字段,您可以决定如何在 MyErrorStatsCombiner 中聚合/保留这些字段。

    我不清楚您是真的想按 IP 分组还是按多个不同的元数据字段分组。如果您想按多个元数据字段进行分组并获取所有这些字段的计数。这可能是一个有用的参考。 GroupBy using multiple data properties。您可以先按 IP 分组,以及它是成功还是错误。不过,我认为您将无法通过同一记录中的总错误和成功获得所需的输出。例如,您可以使用 bigquery 查询轻松完成最后一部分。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2019-01-07
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2014-11-24
      相关资源
      最近更新 更多