【问题标题】:flink aggregate state is huge, how to fixflink聚合状态巨大,如何解决
【发布时间】:2020-01-08 16:52:26
【问题描述】:

我尝试用不同的窗口大小计算流中的数据(窗口大小在steam数据中),所以我使用自定义的WindowAssigner和AggregateFunction,但状态很大(窗口范围从一小时到30天)

在我看来,聚合状态只是存储中间结果

有什么问题吗?

public class ElementProcessingTime extends WindowAssigner<Element, TimeWindow> {
    @Override public Collection<TimeWindow> assignWindows(Element element, long timestamp, WindowAssignerContext context) {
        long slide = Time.seconds(10).toMilliseconds();
        long size = element.getTime() * 60 * 1000;
        timestamp = context.getCurrentProcessingTime();

        List<TimeWindow> windows = new ArrayList<>((int) (size / slide));
        long lastStart = TimeWindow.getWindowStartWithOffset(timestamp, 0, slide);
        for (long start = lastStart; start > timestamp - size; start -= slide) {
            windows.add(new TimeWindow(start, start + size));
        }
        return windows;
    }

    @Override public Trigger<FactorCalDetail, TimeWindow> getDefaultTrigger(StreamExecutionEnvironment env) {
        return ElementTimeTrigger.create();
    }

    @Override public TypeSerializer<TimeWindow> getWindowSerializer(ExecutionConfig executionConfig) {
        return new TimeWindow.Serializer();
    }

    @Override public boolean isEventTime() {
        return false;
    }
}

public class CountAggregate implements AggregateFunction<FactorCalDetail, AggregateResult, AggregateResult> {

    @Override public AggregateResult createAccumulator() {
        AggregateResult result = new AggregateResult();
        result.setResult(0.0);
        return result;
    }

    @Override public AggregateResult add(FactorCalDetail value, AggregateResult accumulator) {
        accumulator.setKey(value.getGroupKey());
        accumulator.addResult();
        accumulator.setTimeSpan(value.getTimeSpan());
        return accumulator;
    }

    @Override public AggregateResult getResult(AggregateResult accumulator) {
        return accumulator;
    }

    @Override public AggregateResult merge(AggregateResult a, AggregateResult b) {
        if (a.getKey().equals(b.getKey())) {
            a.setResult(a.getResult() + b.getResult());
        }
        return a;
    }
}

env.addSource(source)
    .keyBy(Element::getKey)
    .window(new ElementProcessingTime())
    .aggregate(new CountAggregate())
    .addSink(new RedisCustomizeSink(redisProperties));

【问题讨论】:

  • 如果你使用的是keyed windows,可以使用RocksDB state backend来减少堆压力ci.apache.org/projects/flink/flink-docs-stable/ops/state/…
  • 我已经用了RocksDBStateBackend incrementalCheckpointing,但是状态太大了,10000条数据1g左右,受不了(;′⌒`)
  • 您看到内存中的状态快速增长吗?您正在运行多少个节点,每个节点上的 CPU / 内存?
  • 测试环境k8s pod 1 core 2g
  • Checkpoints 之间的最小暂停为 1m,状态大小快速增长

标签: apache-flink flink-streaming


【解决方案1】:

当您分配自定义窗口时,状态大小可能会很快失控。这主要是因为每个窗口都需要保存属于它的所有记录,直到窗口被聚合并最终被驱逐。在您的代码中,您似乎也为每条记录创建了大量窗口。

您没有指定您的用例,但我假设您实际上想要计算每个键在 10 毫秒 bin 大小的给定时间点上有多少事件延伸。如果是这样,那么这不是 windows 的直接用例。

你想做的是:

  1. 将您的活动拆分为更小的活动。
  2. 按 key 和 bin 分组。
  3. 清点你的垃圾箱。

代码中的粗略草图:

input.flatMap(element -> {
        ...
        for (long start = lastStart; start > timestamp - size; start -= slide) {
            emit(new KeyTime(key, start));
        }
    })
    .keyBy(keyTime -> keyTime)
    .count()

您可以在keyBy 之后应用窗口来强制某些输出属性,例如等待几分钟然后输出所有内容并忽略迟到的事件。

注意:KeyTime 是一个简单的 POJO,包含密钥和 bin 时间。

编辑:在您发表评论后,解决方案实际上要简单得多。

env.addSource(source)
    .keyBy(element -> new Tuple2<>(element.getKey(), element.getTime()))
    .count()
    .addSink(new RedisCustomizeSink(redisProperties));

【讨论】:

  • 感谢您的回复,我会尝试这种方式。就我而言,我想在不同的时间跨度内计算密钥。例如此时流中有5个数据: (key1, 10min, data1) , (key1, 10min, data2), (key1, 20min, data3), (key1, 20min, data4), (key1, 20min, data5) ;所以这个点结果是 (key1, 10min, 2), (key1, 20min, 3)
  • 查看我的更新答案以获得更简单的解决方案。顺便说一句,从你的实际用例开始你的问题总是有意义的。
【解决方案2】:

你没有说源是什么,它会有自己的状态持续存在。您也没有说有多少个唯一键。随着唯一键数量的增加,即使每个键的少量状态也会变得巨大。如果问题最终出现在聚合器状态的增长中,您可以尝试将窗口逻辑拆分为一系列两个窗口,一个用于每小时聚合,另一个用于将每小时汇总聚合到您想要的时间范围。

【讨论】:

  • 这似乎更像是一个评论而不是一个答案。
  • 两个windows的方式我不明白,你能解释一下吗
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2014-07-17
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多