Apache Flink 作为分布式流批一体计算框架,性能调优通常从资源、并行度、状态、网络、算子、数据倾斜、Checkpoint、外部系统等多个维度入手。下面系统性地总结常见的 Flink 性能调优方法(以生产实践为主)。
parallelism.default: 4
stream.keyBy(...).map(...).setParallelism(8);
taskmanager.numberOfTaskSlots: 4
taskmanager.memory.process.size: 4096m
stream.filter(...).slotSharingGroup("group1");
backPressuredTimeMsPerSecond| 原因 | 解决 |
|---|---|
| 某算子慢 | 提高并行度 |
| 数据倾斜 | 优化 key |
| 外部系统慢 | 限流 / 批量写 |
| GC | 调 JVM 参数 |
state.backend: rocksdb
state.backend.incremental: true
state.backend.rocksdb.block.cache-size: 64m
state.backend.rocksdb.writebuffer.size: 32m
StateTtlConfig ttl = StateTtlConfig
.newBuilder(Time.days(7))
.build();
execution.checkpointing.interval: 30s
execution.checkpointing.timeout: 10min
execution.checkpointing.mode: EXACTLY_ONCE
execution.checkpointing.enable-async: true
state.backend.incremental: true
taskmanager.network.memory.fraction: 0.1
taskmanager.network.memory.max: 512mb
execution.buffer-timeout: 100ms
keyBy / rebalancekeyBy(key + "-" + random(0,10))
局部聚合 → 打散 → 全局聚合
// 差
map(x -> new Tuple2<>(x, 1))
// 好
reuse object
KafkaSource.builder()
.setProperty("fetch.max.wait.ms", "500")
.setProperty("max.poll.records", "1000")
records-in / records-outcheckpoint durationGC timebackpressure ratio如果你愿意,可以告诉我:
我可以直接给你一份针对性调优方案或参数模板。
免责声明:本站发布的内容(图片、视频和文字)以原创、转载和分享为主,文章观点不代表本网站立场,如果涉及侵权请联系站长邮箱:is@yisu.com进行举报,并提供相关证据,一经查实,将立刻删除涉嫌侵权内容。