MSF 流式分析:CPU 利用率优化 & KPU 成本缩减方案

Amazon Managed Service for Apache Flink · 内存密集型场景调优指南

问题描述

客户在使用 MSF 做流式分析时发现 CPU 利用率持续偏低,业务场景对内存需求更大。由于 KPU 是 CPU/内存/存储的固定捆绑,为获取足够内存而分配了过多 KPU,导致 CPU 闲置、成本浪费。

KPU 资源定义

资源每 KPU说明
CPU1 vCPU单核虚拟 CPU
内存4 GB3 GB 应用 + 1 GB State Store
存储50 GB本地磁盘(State 溢出)
KPU 数量 = ⌈ Parallelism / ParallelismPerKPU ⌉ + 1(编排节点)

KPU 是固定比例捆绑(1:4:50),无法单独调整 CPU 与内存的比例。ParallelismPerKPU 默认为 1,最大为 8。

优化方案

推荐 · 零代码

方案一:提高 ParallelismPerKPU

将每 KPU 承载的 Slot 数从默认的 1 调高到 2~8,使更多 Task 共享同一 KPU 的资源。

效果:保持 Parallelism 不变,直接减少 KPU 数量,CPU 利用率成比例提升。

Parallelism = 8, ParallelismPerKPU = 1 → 8 KPU
Parallelism = 8, ParallelismPerKPU = 2 → 4 KPU(成本 ↓50%)
Parallelism = 8, ParallelismPerKPU = 4 → 2 KPU(成本 ↓75%)

操作路径:MSF 控制台 → Application → Scaling → ParallelismPerKPU → 调高 → Update

适用场景:I/O 密集型、等待外部数据库/Kafka/S3 的场景(CPU 本身不吃紧)

需改代码

方案二:Operator 级别精细化并行度

不让所有 Operator 都使用应用级 Parallelism,而是按资源消耗差异化设置。

Operator 类型建议并行度原因
Source(Kafka/Kinesis)高(匹配 Partition 数)吞吐对齐
Map / Filter(轻量)低(1/4 应用级)CPU 消耗极小
窗口聚合 / Async I/O重状态 / 重 I/O
Sink匹配下游写入能力

稳定比参考:总 Operator 并行度 : Slot = 4:1(资源密集型 2:1~3:1,轻量型可达 10:1)

推荐写法:Runtime Properties 动态配置

// 从 Runtime Properties 读取比例(无需重新编译即可调优) Map<String, Properties> props = KinesisAnalyticsRuntime.getApplicationProperties(); int base = env.getParallelism(); int mapRatio = Integer.parseInt( props.get("OperatorProperties").getProperty("MapRatio", "4")); source.setParallelism(base); // Source 全并行 mapper.setParallelism(base / mapRatio); // Map 1/4 并行 sink.setParallelism(base / 2); // Sink 半并行

首次需改代码植入 Runtime Properties 读取逻辑;后续调优只需在控制台改参数 → Update 即可。

需改代码

方案三:代码级内存优化

降低单任务内存占用,间接减少所需 KPU 总数。

方向做法效果
State Backend使用 RocksDB(State 自动 spill 到磁盘)堆内存大幅释放
State TTL为 Keyed State 设置合理 TTL避免状态无限增长
序列化Avro / Protobuf 替代 Java 序列化内存 & 网络开销降低
窗口AggregateFunction 替代 ProcessWindowFunction增量聚合不缓存全量
数据倾斜检查 KeyBy 分布均匀性避免单 Slot 内存膨胀
高级

方案四:自定义 Metric-Based Scaling

MSF 内置 Autoscaling 仅基于 CPU(≥75% 持续 15 分钟扩、<10% 持续 6 小时缩)。对于 CPU 一直低但内存紧张的场景,它永远不会主动缩容。

建议使用 CloudWatch Alarm + Lambda/Step Functions 实现基于 heapMemoryUtilization、背压、自定义指标的精细缩扩容。

参考:AWS Blog - Enable metric-based and scheduled scaling

极端场景

方案五:迁移到 EMR Flink

如果 CPU:内存比严重不匹配(如需 8GB 内存但仅 0.2 vCPU),MSF 的固定 1:4 比例无法满足。可考虑 EMR 上的 Flink,自由选择内存优化型实例(如 r6g / r7g),精确匹配工作负载。

Trade-off:失去 Serverless 便利性,需自行管理集群;但资源配比完全灵活。

推荐实施步骤

观测当前指标
CloudWatch 中查看 containerCPUUtilization、containerMemoryUtilization、heapMemoryUtilization,确认 CPU 低 + 内存高的模式
调高 ParallelismPerKPU(零代码)
从 1 逐步提升到 2 → 4,每次观察 CPU 是否进入 50-75% 健康区间
调整 Operator 并行度(需改代码)
植入 Runtime Properties 逻辑,对轻量 Operator 降低并行度
优化 State 管理
切换 RocksDB Backend、设置 TTL、优化序列化,减少内存占用
性能测试验证
每步调整后监控延迟和背压,确保处理能力不退化