flink 任务优化系列
flink 任务优化系列
维表更新导致事实表数据全量更新
1、使用支持upsert 模式后计算存储架构,如:flink + paimon
2、事实表只存储维度表id , 查询时在关联
流式聚合优化
MiniBatch 聚合
默认情况下,无界聚合算子是逐条处理输入的记录,即:(1)从状态中读取累加器,(2)累加/撤回记录至累加器,(3)将累加器写回状态,(4)下一条记录将再次从(1)开始处理。这种处理模式可能会增加 StateBackend 开销(尤其是对于 RocksDB StateBackend )。此外,生产中非常常见的数据倾斜会使这个问题恶化,并且容易导致 job 发生反压。
MiniBatch 聚合的核心思想是将一组输入的数据缓存在聚合算子内部的缓冲区中。当输入的数据被触发处理时,每个 key 只需一个操作即可访问状态。这样可以大大减少状态开销并获得更好的吞吐量。但是,这可能会增加一些延迟,因为它会缓冲一些记录而不是立即处理它们。这是吞吐量和延迟之间的权衡。
下图说明了 mini-batch 聚合如何减少状态操作。

任务参数配置:
// instantiate table environment
TableEnvironment tEnv = ...;
// access flink configuration
TableConfig configuration = tEnv.getConfig();
// set low-level key-value options
configuration.set("table.exec.mini-batch.enabled", "true"); // enable mini-batch optimization
configuration.set("table.exec.mini-batch.allow-latency", "5 s"); // use 5 seconds to buffer input records
configuration.set("table.exec.mini-batch.size", "5000"); // the maximum number of records can be buffered by each aggregate operator task
总结:
实时微批处理和离线任务按分钟调度增量计算逻辑一样(实时任务要一直占用计算资源,离线任务会暂时释放计算资源但是任务进程启动耗时较大)
Local-Global 聚合
Local-Global 聚合是为解决数据倾斜问题提出的,通过将一组聚合分为两个阶段,首先在上游进行本地聚合,然后在下游进行全局聚合,类似于 MapReduce 中的 Combine + Reduce 模式
拆分 distinct 聚合
在 distinct 聚合上使用 FILTER 修饰符
参考flink 官网:https://nightlies.apache.org/flink/flink-docs-release-1.18/zh/docs/dev/table/tuning/
任务背压优化
当 Flink 作业正运行在严重的背压下时,Checkpoint 端到端延迟的主要影响因子将会是传递 Checkpoint Barrier 到 所有的算子/子任务的时间。这在 checkpointing process) 的概述中有说明原因。并且可以通过高 alignment time and start delay metrics 观察到。 当这种情况发生并成为一个问题时,有三种方法可以解决这个问题:
1、消除背压源头,通过优化 Flink 作业,通过调整 Flink 或 JVM 参数,抑或是通过扩容。
2、减少 Flink 作业中缓冲在 In-flight 数据的数据量。
3、启用非对齐 Checkpoints。 这些选项并不是互斥的,可以组合在一起。本文档重点介绍后两个选项。
参考flink官网:https://nightlies.apache.org/flink/flink-docs-release-1.18/zh/docs/ops/state/checkpointing_under_backpressure/
优化 Flink 配置
涉及多个方面,包括调整参数以提高性能、资源利用率和容错能力。以下是一些建议的优化措施,你可以根据自己的应用场景和需求进行调整:
1. 调整并行度
任务并行度:增加任务的并行度可以提高吞吐量,减少延迟。但过高的并行度也会增加资源消耗和协调开销。
State Backend 并行度:调整 State Backend 的并行度可以优化状态存储和恢复的性能。
命令行参数配置
java flink run -c com.example.MyJob --state.backend.parallelism 4 my-job.jar
编码方式参数配置
java StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); FsStateBackend stateBackend = new FsStateBackend("hdfs:///path/to/state/directory"); env.setStateBackend(stateBackend); env.setStateBackendParallelism(4); // 设置并行度为 4
-
调整 Checkpoint 配置
Checkpoint 间隔:减小 Checkpoint 间隔可以提高容错能力,但也会增加性能开销。找到一个平衡点,使得 Checkpoint 的开销和容错能力之间达到最优。
Checkpoint 超时时间:合理设置 Checkpoint 超时时间,避免由于超时导致的 Checkpoint 失败。
Checkpoint 模式:根据应用场景选择合适的 Checkpoint 模式(如精确一次性语义或至少一次语义)。
增量 Checkpoint:使用增量 Checkpoint 以减少恢复时的数据量和时间。
-
调整内存配置
TaskManager 内存:根据任务需求调整 TaskManager 的内存配置,包括 JVM Heap 大小和 Off-Heap 内存。
JVM Heap 大小:根据你的作业需求调整 JVM Heap 的大小。Heap 内存主要用于存储对象实例,包括任务状态、中间结果等。如果 Heap 内存不足,可能会导致频繁的垃圾回收,甚至 OutOfMemoryError。可以通过调整 taskmanager.memory.heap.size 来设置 Heap 内存大小。 Off-Heap 内存:Off-Heap 内存主要用于直接缓冲区和某些数据结构,比如网络缓冲区。增加 Off-Heap 内存可以提高 Flink 的吞吐量,减少垃圾回收的影响。可以通过设置 taskmanager.memory.off-heap.size 来配置 Off-Heap 内存大小。网络缓冲区大小:增加网络缓冲区大小可以减少网络拥塞和延迟。
网络缓冲区:Flink 使用网络缓冲区来暂存数据,以便在网络传输过程中进行缓冲。增加网络缓冲区的大小可以减少网络拥塞和延迟。可以通过设置 taskmanager.network.buffers.memory.size 来调整网络缓冲区的大小。 -
调整任务调度
调度策略:根据集群资源情况和任务需求选择合适的调度策略,如公平调度或优先级调度。
Slot 共享:启用 Slot 共享可以让多个任务共享同一个 TaskManager 的资源,提高资源利用率。
-
监控和调优
性能监控:使用 Flink 提供的 Web UI 或其他监控工具实时监控任务的性能指标,如吞吐量、延迟和 Checkpoint 频率等。
日志分析:定期分析 Flink 任务的日志,发现潜在的性能问题和错误,并进行相应的调优。
-
其他优化措施
启用 Watermark:对于时间窗口聚合任务,启用 Watermark 可以处理乱序事件,确保计算的正确性。
优化 State Backend:选择高性能的 State Backend(如 RocksDB)以优化状态存储和恢复的性能。
减少外部系统依赖:尽量减少任务对外部系统的依赖,降低外部系统对任务性能的影响。
在调整 Flink 配置时,建议逐步调整参数,并在每次调整后进行性能测试和监控,以便及时发现和解决性能问题。同时,也要注意保持与 Flink 社区和官方文档的同步,了解最新的优化建议和实践经验。
更多推荐




所有评论(0)