【问题标题】:Bad Perf with Sliding Windows in FlinkFlink 中滑动窗口的性能不佳
【发布时间】:2017-09-13 15:24:43
【问题描述】:

我使用此代码执行我的测试 (Flink Quick Start):

 val text = env.socketTextStream("localhost", port, '\n')

    // parse the data, group it, window it, and aggregate the counts
    val windowCounts = text
        .flatMap { w => w.split("\\s") }
        .map { w => WordWithCount(w, 1) }
        .keyBy("word")
        .timeWindow(Time.minute(15))
        .sum("count")

使用此代码,我有超过 65 000 次输入/秒

如果我改变了

timeWindow(Time.minute(15))

timeWindow(Time.minutes(15), Time.seconds(1))

我的输入/秒数少于 2 500

有什么方法可以让滑动窗口获得更好的性能?

【问题讨论】:

    标签: scala apache-flink flink-streaming


    【解决方案1】:

    对于 15 分钟滚动窗口,每个传入事件都分配到一个窗口,而对于 15 分钟滑动窗口和一秒幻灯片,每个传入事件被复制到 15 * 60 = 900 个窗口中。这显然会对性能产生影响。

    根据您的应用程序要求,您可以通过使用 ProcessFunction 或通过实现自定义窗口逻辑来以更少的开销计算您需要的内容。例如,您可以预先聚合成 900 个一秒窗口,然后有第二层窗口,通过减去过期秒对总数的贡献并添加最近一秒的值来增量调整 15 分钟的结果。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2017-09-09
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2014-03-28
      • 1970-01-01
      • 2014-07-05
      • 1970-01-01
      相关资源
      最近更新 更多