ARTICLE DETAIL

资讯详情

深耕网站建设、视觉设计与SEO优化的一线实战洞察。

Flink基础之Flink并行度和Slot深度剖析:资源到底怎么算

Flink基础之Flink并行度和Slot深度剖析:资源到底怎么算 摘要讲透 Flink 并行度与 Slot 的完整机制并行度与 Slot 的本质区别与映射关系、算子链与 Slot Sharing 如何决定资源需求、并行度四级配置优先级以及资源规划三步法并给出并行度调整的五个真实踩坑点。关键词Flink、并行度、Parallelism、Slot、Slot Sharing、算子链、配置优先级、资源规划、数据倾斜、keyBy很多 Flink 新手和一部分老手都在这两个词上栽过跟头并行度和Slot。调并行度时随手写个-p 16发现集群上任务还是排队或者以为「并行度开得越大越快」结果数据倾斜把吞吐拉崩。根源在于没搞清楚一个核心问题并行度是逻辑概念Slot 是物理概念二者之间隔着一层「算子链 Slot Sharing」的映射。这一篇把它彻底讲透最后给出一套可以照着算的资源规划方法。一、并行度与 Slot逻辑与物理的分工先分清两个概念的本质并行度Parallelism逻辑层面的概念指一个算子被拆成几个并行实例。map并行度 3就是 3 个 map 实例同时跑每个实例处理一部分数据分区。Slot物理层面的概念TaskManager 内部划分的资源单元。taskmanager.numberOfTaskSlots决定每个 TM 有几个 Slot集群总 Slot 数 Σ(各 TM 的 Slot 数)。它们的关系是一条硬约束作业的实际并行度 ≤ 集群可用 Slot 数。Slot 不够任务就只能停在 SCHEDULED 状态排队。还有个必须澄清的误区Slot 是内存划分不是 CPU 隔离。同一个 TM 的多个 Slot 共享同一批 CPU 核心Slot 只约束并发任务数和内存配额。要 CPU 隔离靠 YARN/K8s 的容器粒度不靠 Slot。二、任务到 Slot 的映射算子链与槽位共享并行实例不会一个个单独落到 Slot 上中间有两层机制在起作用算子链OperatorChain并行度相同 数据 one-to-one 传播的算子合并成一个 Task在同一线程内串行执行。Source map filter可能只是一个 Task。Slot Sharing槽位共享默认情况下同一作业的不同 Task并行度匹配可以共享同一个 Slot。这两层叠加的效果很反直觉但非常重要一个作业需要的 Slot 数 ≈ 作业中最大的并行度而不是所有算子并行度之和。比如作业Source(并行度2) → keyBy → 聚合(并行度4) → Sink(并行度4)不共享需要 2 4 4 10 个 Slot默认共享只需要 max(2,4,4) 4 个 Slot并行度小的链「补位」进共享 Slot。这是 Flink 资源利用率高的核心原因也是理解「为什么我开了 16 并行度集群却只用了 4 个 Slot」的钥匙——要么被算子级配置覆盖了要么共享组内并行度根本没到 16。三、并行度的四级配置优先级并行度可以在四个层面配置优先级从高到低算子级最高map(...).setParallelism(8)只影响单个算子调优时最常用。执行环境级env.setParallelism(4)作为作业内所有算子的默认值。提交参数flink run -p 6不改代码快速调整作业默认并行度A/B 验证常用。配置文件parallelism.default全局兜底优先级最低。一个特殊点Source 的并行度不完全受这套优先级控制。Kafka Source 的并行度默认取min(分区数, parallelism)文件 Source 可以显式指定。这就是为什么有时候你设了并行度 8Source 却只有 4 个实例——topic 只有 4 个分区。排查技巧Web UI 上每个算子旁都会显示实际生效的并行度与预期不符时按四级逐层往下查一定是某一级配置压住了。四、Slot Sharing 深入资源到底怎么算理解了默认共享再看自定义分组和资源规划4.1 默认共享 vs 分组隔离默认共享所有算子属于default组Slot 需求 ≈ 最大并行度。省资源但重任务聚集在同一 Slot 时会互相挤占。自定义分组算子.slotSharingGroup(g1)把任务分组同组共享 Slot、不同组互不干扰。比如把重聚合单独分一组避免它和 Source 抢同一个 Slot 的线程资源。分组是「用资源换隔离」默认共享省资源但可能互相拖累分组隔离保性能但多占 Slot。Slot 需求 Σ(各组内最大并行度)。4.2 资源规划三步法确定各算子并行度考虑数据量、状态规模、key 分布。数据量小就别开高并行度key 数量少时并行度高只会让大部分子任务空转。估算 Slot 需求默认共享取最大并行度有分组时按组内最大并行度求和。反推 TM 数量TM 数 × 每 TM Slot 数 ≥ Slot 需求并留 20% 余量应对流量波动和单 TM 故障时的任务迁移。示例Slot 需求 8、每 TM 4 Slot → 2 个 TM如果要求单 TM 故障不影响作业需 1 冗余则要 3 个 TM——这也是生产环境常见的「多一台备着」的由来。五、并行度调整的五个真实踩坑并行度 总 Slot 数作业永远起不来。YARN 上表现为作业一直 ACCEPTED/RUNNING 但不推进Web UI 里任务全是 SCHEDULED。先看集群还剩多少 Slot再定并行度。盲开高并行度吞吐反而下降。数据倾斜时一个 key 的数据全压在一个子任务上并行度再高其他实例都在空转而且并行度翻倍网络 shuffle 的链接数也翻倍并行度²序列化开销同步上涨。先解决倾斜再谈并行度。keyBy 后并行度变化触发状态重分布。有状态算子的并行度调整无论是改配置还是从 Savepoint 恢复时改动都会触发 key group 重分布作业恢复时有一段重算开销。无界流改并行度必须基于 Savepoint 恢复直接改配置重启会丢状态语义。固定 key 让并行度形同虚设。keyBy(v - total)把所有数据压到单并行度其他实例空转——这是「并行度 16 但实际只有 1 个在干活」的最常见原因。key 设计要保证分布均匀。忽略 Slot Sharing 直接按并行度之和申请资源。不知道默认共享机制的人会按「所有算子并行度相加」去配 TM 数量白白申请两三倍的资源。按最大并行度估算能省一大半。并行度与 Slot 的关系可以浓缩成一句话并行度决定「逻辑上要几个实例」Slot 决定「物理上能放几个任务」中间由算子链和 Slot Sharing 桥接——所以资源需求 ≈ 最大并行度而非并行度之和。把「四级优先级」「默认共享省资源」「key 分布决定真实吞吐」这三件事记牢Flink 的资源规划与并行度调优就都有了依据。算子链和 Slot Sharing 桥接**——所以资源需求 ≈ 最大并行度而非并行度之和。把「四级优先级」「默认共享省资源」「key 分布决定真实吞吐」这三件事记牢Flink 的资源规划与并行度调优就都有了依据。
返回列表