Flink如何统计KeyedStream ProcessFunction并行实例的指标
Flink KeyedStream下ProcessFunction并行实例维度指标采集方案
Flink原生支持该场景下的并行实例维度指标采集,无需拆分多作业,基于内置指标体系即可实现,具体方案如下:
核心逻辑
KeyedStream下ProcessFunction的普通类成员变量是并行实例级别的共享变量,不会随key的上下文切换重置,且算子单线程处理所有流经的key,不存在并发竞争问题,完全可以用来统计整个实例的全局状态,不需要依赖Keyed State做单key维度的统计再聚合。
落地步骤
- 初始化实例级指标与统计变量
在open()初始化方法中定义非State的普通成员变量,用来记录当前并行实例的全局统计值,同时通过Flink RuntimeContext提供的MetricGroup注册对应监控指标:// 注意:以下是普通类成员变量,不属于Keyed State,为当前并行实例全局共享 private long instanceTotalBufferedBytes = 0L; private int instanceActiveTopicCount = 0; // 额外维护每个key的缓存大小映射,方便全局flush时快速定位大批次key private Map<String, Long> keyBufferedBytesMap = new HashMap<>(); @Override public void open(Configuration parameters) { MetricGroup opMetricGroup = getRuntimeContext().getMetricGroup() .addGroup("kafka_to_gp_sync") .addGroup("batch_process"); // 注册实例总缓存字节数指标 opMetricGroup.gauge("instance_buffered_bytes", () -> instanceTotalBufferedBytes); // 注册实例当前正在攒批的主题数指标 opMetricGroup.gauge("instance_active_topic_count", () -> instanceActiveTopicCount); } - 消息处理时同步更新全局统计
每处理一条消息,除了将消息写入当前key对应的Keyed State,同步计算消息字节数:- 给
instanceTotalBufferedBytes累加当前消息字节数 - 更新
keyBufferedBytesMap中当前主题的缓存字节数 - 如果是该主题进入当前批次的第一条消息,给
instanceActiveTopicCount加1
- 给
- 批次flush时同步扣减全局统计
无论是字节阈值触发、等待时长定时器触发的flush操作,在将某主题的批次数据写入Greenplum并清空对应Keyed State后:- 从
instanceTotalBufferedBytes中扣除该主题本次flush的总字节数 - 从
keyBufferedBytesMap中移除该主题的统计记录 - 给
instanceActiveTopicCount减1
- 从
OOM问题配套优化
基于实例级的全局统计,可直接增加实例级内存保护机制:当instanceTotalBufferedBytes达到预设的单实例内存警戒阈值(可设置为算子托管内存的60%~70%),强制flush当前实例中缓存字节数最大、或首条消息等待时间最长的主题批次,避免多主题同时攒批导致总内存超出容器限制触发OOM。
方案对比说明
不需要采用按主题拆分独立作业的方案,该方案会带来数百个作业的运维成本,且资源利用率极低;也不需要将攒批触发逻辑绑定到checkpoint阶段,该方式会让攒批节奏完全受checkpoint间隔限制,无法灵活满足字节阈值、等待时长阈值任意触发的业务需求。
内容的提问来源于stack exchange,提问作者Alexey Ryabov
相关产品推荐
相关产品推荐

