SeaTunnel Zeta 引擎资源隔离实战:用节点 tag 与 tag_filter 精确控制作业调度
【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel
SeaTunnel 的 Zeta(SeaTunnel 自研引擎)支持为每个 Worker 节点添加tag,并在作业配置文件中通过tag_filter声明式地选择作业要运行的节点,从而实现多团队、多业务线共享同一集群时的资源隔离。本文完整讲解「节点打标签 → 作业按标签过滤 → 运行中动态更新标签」的全流程,并结合仓库源码剖析tag_filter的匹配规则、候选节点为空时的异常处理机制,以及 E2E 测试对该功能的验证方式。读完后,你能够在一套共享的 Zeta 集群上为不同团队圈定专属 Worker 池,并能通过 REST API 在不重启节点的情况下调整其归属。
1. 功能定位与整体原理
资源隔离要解决的问题是:同一套 SeaTunnel 集群往往承载多个团队(或不同优先级、不同数据区域)的作业,如果不加约束,作业会被调度到集群中任意一个有空闲 Slot 的 Worker 上,业务之间会互相争抢 CPU 和内存。Zeta 引擎提供的隔离机制由两部分组成:
- 节点侧(供给侧):通过 Hazelcast 的
member-attributes给每个集群成员打上一组 key-value 形式的tag(例如group=platform, team=team1)。节点注册到 ResourceManager 时,这些 attributes 会随WorkerProfile一起上报,成为该节点的"标签画像"。 - 作业侧(需求侧):在作业的
env配置块中声明tag_filter,ResourceManager 在为该作业申请 Slot 前,先用tag_filter过滤出候选 Worker 集合,再在候选集合内按分配策略挑选 Slot。
从源码结构看,tag_filter被定义为一个Map<String, String>类型的 env 选项,定义在 EnvCommonOptions.java:
public static Option<Map<String, String>> NODE_TAG_FILTER = Options.key("tag_filter") .mapType() .noDefaultValue() .withDescription("Define the worker where the job runs by tag");该选项没有默认值,且与custom_parameters等 env 选项一起注册进 EnvOptionRule 中,保证它能被作业配置的env { ... }块合法解析。
需要说明适用前提:tag_filter调度过滤是 Zeta 引擎的能力。仓库中的 E2E 测试 ResourceIsolationIT 就显式标注了@DisabledOnContainer(value = {}, type = {EngineType.SPARK, EngineType.FLINK}, disabledReason = "only work on Zeta"),说明该特性仅在 Zeta 引擎下生效。
2. 第一步:在 hazelcast.yaml 中为节点打 tag
以发行包自带的 config/hazelcast.yaml 为基础更新配置,在hazelcast根节点下增加member-attributes段:
hazelcast: cluster-name: seatunnel network: rest-api: enabled: true endpoint-groups: CLUSTER_WRITE: enabled: true DATA: enabled: true join: tcp-ip: enabled: true member-list: - localhost port: auto-increment: false port: 5801 properties: hazelcast.invocation.max.retry.count: 20 hazelcast.tcp.join.port.try.count: 30 hazelcast.logging.type: log4j2 hazelcast.operation.generic.thread.count: 50 member-attributes: group: type: string value: platform team: type: string value: team1在这个配置中,我们通过member-attributes设置了group=platform、team=team1两个tag。member-attributes下的每个属性由type(取值如string)和value组成,属性名即 tag 的 key,value即 tag 的 value。
配置要点:
- 该文件对每个需要打标签的节点生效,即同一集群中不同机器上的
hazelcast.yaml可以各自配置不同的member-attributes,从而把集群划分成若干"标签域"; - 每个节点的 tag 集合是相互独立的,例如节点 A 可以配置
team=team1,节点 B 配置team=team2,不配置member-attributes的节点则没有任何标签; - 修改
member-attributes需要节点以新配置启动后生效;对已在运行的节点,见第 5 节的 REST API 动态更新方式。
3. 第二步:在作业配置中使用 tag_filter
在作业配置文件的env块中声明tag_filter,把作业"锚定"到匹配标签的节点上:
env { parallelism = 1 job.mode = "BATCH" tag_filter { group = "platform" team = "team1" } } source { FakeSource { plugin_output = "fake" parallelism = 1 schema = { fields { name = "string" } } } } transform { } sink { console { plugin_input="fake" } }上面的示例与仓库 E2E 用例 fakesource_to_console.conf 的结构一致(E2E 用例的 schema 多定义了id、age字段,核心tag_filter写法相同)。
3.1 匹配语义(重点)
- 不配置
tag_filter:作业会从所有已注册节点中随机选择节点来运行,不做任何标签约束; - 配置了多个过滤条件:采用"与"(AND)语义,要求节点标签中每个 key 都存在且 value 完全相等,任一条件不满足该节点即被排除;
- 没有任何节点匹配:抛出
NoEnoughResourceException,作业无法启动。
3.2 源码级过滤逻辑
过滤逻辑集中在 AbstractResourceManager.filterWorkerByTag:
private ConcurrentMap<Address, WorkerProfile> filterWorkerByTag(Map<String, String> tagFilter) { if (tagFilter == null || tagFilter.isEmpty()) { return registerWorker; // 未配置 tag_filter:候选集为全部已注册 worker } return registerWorker.entrySet().stream() .filter(e -> { Map<String, String> workerAttr = e.getValue().getAttributes(); if (workerAttr == null || workerAttr.isEmpty()) { return false; // 节点没有任何 tag,直接不匹配 } boolean match = true; for (Map.Entry<String, String> entry : tagFilter.entrySet()) { if (!workerAttr.containsKey(entry.getKey()) || !workerAttr.get(entry.getKey()).equals(entry.getValue())) { return false; // key 缺失或 value 不等,均判为不匹配 } } return match; }) .collect(Collectors.toConcurrentMap(Map.Entry::getKey, Map.Entry::getValue)); }从这段实现可以确认三个事实:
tagFilter为null或空 Map 时直接返回全部已注册 Worker,对应"未配置则随机选择"的行为(后续 Slot 选择再交由分配策略,默认RandomStrategy);- 节点属性为空的 Worker永远不会匹配任何非空
tag_filter; - 条件是"key 必须存在且 value 相等"的严格全量匹配,不支持模糊或前缀匹配,也不支持"任一匹配"(OR)语义。
3.3 候选集为空时的异常路径
在 AbstractResourceManager.applyResources 中,过滤结果会立即做空判断:
ConcurrentMap<Address, WorkerProfile> matchedWorker = filterWorkerByTag(tagFilter); if (matchedWorker.isEmpty()) { log.error("No matched worker with tag filter {}.", tagFilter); throw new NoEnoughResourceException(); }即"有 tag_filter 但没有一个节点满足"会直接抛出 NoEnoughResourceException,与"节点匹配但 Slot 不足"共用同一异常类型。作业在提交阶段读取该选项的位置是 PhysicalPlanGenerator,它从作业配置的 env options 中取出NODE_TAG_FILTER.key()即"tag_filter"的 Map 值,随每个 Pipeline 的资源申请向下传递。
3.4 E2E 测试验证
仓库的 E2E 套件对两种路径都做了断言(见 ResourceIsolationIT):
testTagMatch:执行带tag_filter { group = "platform", team = "team1" }的作业(对应 fakesource_to_console.conf),断言进程退出码为 0;testTagNotMatch:执行 fakesource_to_console_tag_not_match.conf(配置了集群中不存在的标签组合),断言退出码非 0,且 stderr 中包含org.apache.seatunnel.engine.server.resourcemanager.NoEnoughResourceException。
4. 典型使用场景
标签机制本质上是给 Worker 池做"软分区",常见用法:
- 多租户/多团队隔离:
tag_filter { tenant = "a" }将租户 A 的作业限制在打了tenant=a标签的节点上,租户间不互相侵占 Slot; - 数据局部性:按可用区/机房打标(如
zone = "us-west-1"),把 ETL 作业调度到与数据源同区域的节点,降低跨机房流量; - 资源专业化:为大内存或特定硬件的节点打
resource = "bigmem"之类的标签,让大状态作业优先落位。
这些用法都是同一套"key-value 全量匹配"语义的组合,无需引入额外组件。
5. 可选:运行时通过 REST API 更新节点 tags
节点标签不必写死在hazelcast.yaml中,运行中的节点也可以动态增删标签。该接口由 Zeta server 内嵌的 HTTP 服务提供(前提是seatunnel.engine.http.enable-http = true或enable-https = true,默认发行包配置已开启并监听 8080),详见 rest-api-v2.md。
5.1 更新节点 tags
对指定节点的ip:port调用POST /update-tags,请求体是一个 Map 对象,表示要写入该节点的新 tags:
POST /update-tags成功时返回:
{ "message": "update node tags done." }5.2 清除节点 tags
请求体传空 Map 对象,表示清除当前节点的全部 tags。
由于更新是针对单节点的,需要使用目标节点的ip:port定位。这意味着集群上线后,可以在不重启任何 Worker 的前提下,把某台机器从team1池迁移到team2池,或把一台节点从任何标签池中"摘除"(清除标签后它将不再匹配任何非空tag_filter,但仍可被无过滤条件的作业随机选中)。
5.3 验证标签是否生效
更新后可以通过以下只读接口核对:
GET /resource/workers:返回每个已注册 Worker 的资源快照,其中tags字段即节点当前标签(无标签的节点返回空对象{});GET /overview?tag1=value1&tag2=value2:按标签过滤后返回满足条件的节点数(works)及其 Slot 汇总,可用于快速确认某个标签域内有多少 Worker。注意runningJobs等作业级指标是集群级别的,不受标签过滤影响。
6. 常见问题与排查
- 提交作业时立即报
NoEnoughResourceException:优先怀疑tag_filter与节点标签拼写不一致(key 或 value 严格区分、全量相等),或对应节点尚未以新member-attributes重启。可通过GET /resource/workers查看各节点实际tags后比对。 - 修改了 hazelcast.yaml 但过滤仍不生效:
member-attributes随成员注册进入 WorkerProfile,节点需以新配置加入集群;对已在运行的节点请使用第 5 节的POST /update-tags。 - 误以为
tag_filter能跨引擎使用:该过滤逻辑位于 Zeta 引擎的 ResourceManager 中(Flink/Spark 引擎下由对应引擎自身负责资源调度),E2E 测试也明确只在 Zeta 容器下运行。 - 标签与 Slot 分配策略的关系:
tag_filter只负责圈定候选节点集合;在候选集合内部具体选哪个节点、哪个 Slot,仍由 slot-allocate-strategy(RANDOM/SLOT_RATIO/SYSTEM_LOAD)决定,两者正交、可叠加使用。
7. 小结
Zeta 引擎的资源隔离围绕「Hazelcastmember-attributes打 tag + 作业env { tag_filter }过滤」展开:节点侧的标签成为WorkerProfile的一部分,作业侧的tag_filter在filterWorkerByTag中做 key 存在且 value 相等的严格 AND 匹配,匹配不到任何节点即抛出NoEnoughResourceException。整套机制无需独立组件、配置即可生效,并可配合POST /update-tags在运行时动态调整节点归属,适合在同一共享集群中实现团队级、区域级的资源隔离。
【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考