bg9.yarn/k8s + dag
大数据里「调度」这个词被严重复用,实际指两个不同层次、彼此正交的问题。动笔前先把它们钉死,后面就不会混。
- (A) 资源调度(resource scheduling):集群一共就这么多 CPU/内存,几十上百个作业同时想跑,谁此刻能拿到多少资源、在哪台机器上跑?——这是「分算力」。
- (B) 工作流编排(workflow orchestration,也叫任务编排 / 作业调度):一条数据管道有上百个任务,彼此有先后依赖、要定时跑、会失败要重试、还要补历史数据,谁该在什么时候被触发、它跑之前谁必须先跑完?——这是「定顺序」。
一句话区分(记住这句,全篇都在展开它):
资源调度回答「给多少算力、放哪跑」;工作流编排回答「什么时候触发、依赖谁先跑完」。前者给算力,后者定顺序。
二者的关系是上下配合,不是二选一:
flowchart TB
O["工作流编排层(定顺序)<br/>Airflow / DolphinScheduler<br/>按 定时+依赖 决定「现在该跑谁」"]
R["资源调度层(给算力)<br/>YARN RM / K8s kube-scheduler<br/>决定「给多少容器、放哪台机器」"]
W["Driver + Executor 真正开始算<br/>(05 篇的进程组)"]
O -->|"到点且上游就绪,把作业提交上去"| R
R -->|"分配容器 / Pod"| W
W -->|"结果落表,通知编排层"| O2["触发下游任务……"]
编排器决定「现在该跑 task X 了」,然后把 X 当成一个 Spark 作业提交给 YARN 或 K8s,由资源调度器给它容器。编排器不管一台机器有几核,资源调度器不管任务之间谁依赖谁。 下面分别展开。
OS 调度的粒度是 CPU 时间片(毫秒级、频繁抢占),分配的是「物理核心的时间」;集群资源调度的粒度是容器(Container / Pod)——一次分配一批 <CPU + 内存> 额度,生命周期是分钟到小时级,分配的是「资源配额」而非「核心的时间片」。所以别把两者的抢占频率、粒度画等号。
YARN 是怎么被逼出来的
- 旧方案(Hadoop 1.x):只有一个
JobTracker,它同时干两件事——管理全集群资源 + 调度和监控每一个 MapReduce 任务。痛点:① 一个进程既分资源又盯任务,是单点也是瓶颈,社区里普遍认为扩展性到约 4000 节点就吃力(这是当年被反复引用的经验数字,具体上限随硬件和版本浮动,可查 Hadoop 官方 YARN 设计文档核实);② 整个集群只能跑 MapReduce,Spark/Flink 这些新框架想共享同一个集群没门。 - YARN 解决什么(Hadoop 2.0,2012–2013 引入):把「资源管理」和「作业的生命周期管理」拆开。资源管理下沉成一个通用层(ResourceManager),每个作业自带一个「作业管家」(ApplicationMaster)管自己的生老病死。这样一个集群能同时跑 MR、Spark、Flink、Tez,共享资源。YARN = Yet Another Resource Negotiator。
- 代价:架构复杂度上升(多了 per-app 的 AM 角色);资源模型偏「粗粒度、内存中心」,一个 Container 分出去就是一整份
<内存, vcores>,不像 OS 那样细到时间片;AM 本身也要占一个 Container。
YARN 是什么:四个角色
flowchart TB
RM["ResourceManager(RM)<br/>全局资源仲裁 = Scheduler + ApplicationsManager"]
subgraph 机器A
NM1["NodeManager A<br/>管本机资源、上报心跳"]
AM["ApplicationMaster<br/>本作业的管家"]
end
subgraph 机器B
NM2["NodeManager B"]
C1["Container<br/>一份 内存+vcores 额度<br/>里面跑 Executor"]
end
RM --- NM1
RM --- NM2
NM1 --- AM
NM2 --- C1
AM -.管理.-> C1
- ResourceManager(RM):全集群唯一的资源大脑。内部再分两块:Scheduler(纯粹按策略把资源分给谁,不负责监控任务死活)和 ApplicationsManager(受理作业提交、为每个作业启动第一个 Container 来跑它的 AM)。
- NodeManager(NM):每台机器一个。管理本机的 CPU/内存额度,定期向 RM 发心跳汇报,负责按命令启动和监控 Container。
- ApplicationMaster(AM):每个作业一个,作业结束就退出。它向 RM 申请资源,拿到 Container 后联系 NM 把任务进程拉起来,并跟踪任务进度。它是「per-application」的——这正是 YARN 相比旧 JobTracker 的关键解耦点:作业级的调度逻辑从中心搬到了每个作业自己身上。
- Container:资源的抽象单位 = 某台机器上的一份
<内存, vcores(虚拟核)>额度。所有实际计算都跑在 Container 里——包括 AM 自己。这个词记牢,下面全链路都靠它。
RM 里的 Scheduler 按哪种策略分?Apache Hadoop 内置三种(术语一致:下文统称「调度器」):
- FIFO Scheduler:先到先得,排一条队。一个大作业进来能把后面全饿死,基本只用于测试。
- Capacity Scheduler(容量调度器):把集群按队列切成若干份容量(如
dept-a队列 60%、dept-b40%),队列之间硬隔离,队列内部再按 FIFO/优先级排。适合多租户按部门分账。这是 Apache Hadoop 当前版本的默认调度器——yarn-default.xml里yarn.resourcemanager.scheduler.class默认值就是CapacityScheduler(我在 Hadoop 3.x 官方yarn-default.xml核实过;注意历史上 CDH 发行版默认用的是 Fair Scheduler,别一概而论)。 - Fair Scheduler(公平调度器):让所有运行中的作业动态公平共享资源——一个作业单独跑时可用满集群,第二个进来时两者动态平分。可配最小保障份额(min share)。适合交互式、作业数量波动大的场景。
自问:Capacity 和 Fair 到底差在哪? 自答:差在「隔离」还是「共享」的取向。Capacity 先按队列切死容量再在队列内排队,强调租户间隔离与容量保证;Fair 强调让活跃作业尽量均分、空闲资源随时被借走。代价对称:Capacity 可能出现「A 队列闲着、B 队列排队却借不满」(除非开弹性借用),Fair 则可能让某个作业的资源被后来者不断稀释、跑得忽快忽慢。
全链路:一个 Spark 作业如何向 YARN 要到资源(兑现 05 篇的悬念)
这是本篇的核心链路。05 篇问「Driver/Executor 的资源从哪来」,答案就是这条时序:
sequenceDiagram
participant C as Client_spark_submit
participant RM as ResourceManager
participant NM1 as NodeManagerA
participant AM as ApplicationMaster_内含Driver
participant NM2 as NodeManagerB
participant EX as Executor
C->>RM: 1. 提交应用(指定队列、AM/Executor 资源需求)
RM->>NM1: 2. 挑一台 NM,分配「首个 Container」
NM1->>AM: 3. 在该 Container 里启动 ApplicationMaster
Note over AM: cluster 模式下 Driver 就跑在 AM 里
AM->>RM: 4. 注册,并按 --num-executors/--executor-cores/-memory 申请 Executor 容器
RM-->>AM: 5. 心跳返回「已分配的 Container」(位于 NM B)
AM->>NM2: 6. 拿着授权去请求 NM B 启动 Executor
NM2->>EX: 7. 启动 Executor 进程
EX->>AM: 8. Executor 反向注册到 Driver
AM->>EX: 9. Driver 把 Task 下发给 Executor
EX-->>AM: 10. 汇报进度/结果
AM->>RM: 11. 作业完成,注销并释放所有 Container
对应的提交命令(每行注释;模拟输入 + 结果):
# 向 YARN 提交一个 Spark 作业
spark-submit \
--master yarn \ # 资源来自 YARN 集群(而非本地/standalone)
--deploy-mode cluster \ # cluster 模式:Driver 跑在集群里的 AM 中
--num-executors 10 \ # 要 10 个 Executor 容器
--executor-cores 4 \ # 每个 Executor 4 个 vcore
--executor-memory 8g \ # 每个 Executor 8G 内存
--queue dept-a \ # 提交到 Capacity 调度器的 dept-a 队列
app.py
模拟运行时的关键日志(节选):
INFO Client: Submitting application application_1706... to ResourceManager
INFO Client: Application report for application_1706... (state: ACCEPTED) # 已受理,AM 启动中
INFO Client: Application report for application_1706... (state: RUNNING) # AM 起来了,开始要 Executor
INFO YarnAllocator: Will request 10 executor containers, each with 4 core(s), 9216 MB ...
INFO YarnAllocator: Launching executor on host nodeB-01 # 拿到容器,去 NM 上拉 Executor
一句话解释结果为何如此:注意 9216 MB 而不是你写的 8192(8g)——YARN 会在 Executor 内存之上叠加一份 spark.executor.memoryOverhead(堆外/JVM 开销,默认约为 max(384MB, 0.1×executor内存),随 Spark 版本可能微调,以对应版本官方配置文档为准),所以真正向 YARN 申请的 Container 内存是「你要的 + overhead」。这解释了一个高频困惑:明明只配了 8G,YARN 上看到的容器却更大。
client 模式和 cluster 模式的区别(承接 05 篇):
--deploy-mode cluster时 Driver 在 AM 容器里,提交完就能断开客户端;--deploy-mode client时 Driver 留在你的提交机上,AM 退化成一个只负责申请 Executor 的「跑腿」(ExecutorLauncher),日志直接打在你终端上——所以交互式调试(spark-shell / notebook)用 client 模式,线上生产用 cluster 模式。
YARN 是「大数据专用」的资源层。可是企业里通常还有另一套资源系统在跑在线微服务——K8s。两套资源池各管各的,离线集群半夜忙、白天闲,在线集群反过来,谁也借不到谁的资源,利用率低。
- 解决什么:用 K8s 统一资源池,让离线批作业和在线服务混部在同一批机器上错峰共享;同时配合存算分离(数据在对象存储,计算随起随销)拿到弹性伸缩——高峰拉起几百个计算 Pod,跑完立刻还回去。
- 代价:K8s 原生调度器天生是为「长期运行的在线服务」设计的,面对「一个作业要么全给要么别给」的批处理语义会水土不服(下面 3.3 展开);而且大数据生态在 YARN 上积累了十几年的队列/公平调度经验,迁到 K8s 要重新补齐。
- 成熟度:Spark 官方对 K8s 的支持从 Spark 2.3(2018)实验性起步,到 Spark 3.1.1(2021)正式宣布 Generally Available(GA)(对应 JIRA SPARK-33005,我在 Spark 3.1.1 官方 release notes 核实过)。这是「敢用于生产」的分水岭。
「决定这个 Pod 放到哪台 Node」,就是 kube-scheduler。它干的正是 YARN Scheduler 干的事:在有限的 Node 资源里,给待运行的单元找位置。两个概念直接对上:
requests(请求量)= 调度时的「预定量」。kube-scheduler 按requests做装箱(bin-packing):保证一台 Node 上所有 Pod 的requests之和 ≤ Node 可分配量。它只看 requests 决定放不放得下。limits(上限)= 运行时的「天花板」,由内核 cgroups 强制:CPU 用超limit会被 限流(throttle)(降速,不杀),内存用超limit会被 OOMKill(直接杀)。
和 YARN 的对应(近似,便于迁移直觉):
| 概念 | YARN | K8s |
|---|---|---|
| 资源单元 | Container | Pod |
| 调度依据的「要多少」 | Container 的 <内存, vcores> |
Pod 的 resources.requests |
| 内存超限的下场 | NM 监控到超用,杀掉 Container | 内核 OOMKill 该容器 |
| 谁来放置 | RM 的 Scheduler | kube-scheduler |
kube-scheduler 内部两阶段(对照 OS 调度直觉):先 Filtering(过滤)——筛出「放得下」的 Node(资源够、亲和性/污点满足),类似 OS 从所有线程里筛出「可运行」的那批;再 Scoring(打分)——在候选 Node 里挑最优的(如最空闲的),类似 OS 按优先级排序选一个。这与 YARN「先看队列有没有容量、再挑机器」是同构的。
Spark on K8s 的形态:没有 RM/NM/AM 这套,而是 spark-submit 直接对 K8s API Server 说话——先起一个 Driver Pod,Driver 再向 API Server 申请起若干 Executor Pod,调度全交给 kube-scheduler。数据在对象存储(如 S3、火山的 TOS),Pod 用完即销毁。
为什么 K8s 原生调度器跑批会「资源死锁」
自问:既然 kube-scheduler 会自动放 Pod,为什么大数据圈还要额外装 Volcano、Apache YuniKorn 这类批调度器?
自答:因为 kube-scheduler 是逐个 Pod 独立调度的。一个 Spark 作业要 100 个 Executor,如果集群只够放 60 个,原生调度器会先把这 60 个放上去占着资源,剩下 40 个 pending——而这 60 个占着资源又干不成整活,别的作业也进不来,形成资源死锁 / 饥饿。批作业需要的是 gang scheduling(组调度):要么 100 个一起分配,要么一个都不分。这正是 Volcano、Apache YuniKorn 补的能力(gang scheduling + 队列 + 公平调度)。
现在单个作业能拿到资源跑起来了。但真实数仓不是一个作业,而是几百上千张表按 01 篇的分层层层依赖:ods_order → dwd_order → dws_user_order → ads_gmv_report。早期怎么管?用 Linux crontab 定时 + shell 脚本串起来。痛点很快暴露:
- 依赖靠猜:下游任务不知道上游好没好,只能「估上游 3 点跑完,我排 4 点」,靠
sleep和错峰硬扛——上游一延迟,下游就读到半成品数据。 - 失败没兜底:脚本挂了没人知道,没有自动重试、没有告警。
- 补历史数据是灾难:要重算上个月的数据,只能手改脚本里的日期,跑一天改一次。
- 依赖关系没人看得清:几百个任务的先后关系散在几十个 crontab 里,全靠人脑记。
工作流编排器就是来解决这四点的:把任务组织成 DAG,由它统一按依赖触发、定时、重试、补数。
DAG(有向无环图):节点 = 任务,有向边 = 依赖(A → B 表示 B 依赖 A 的产出);无环意味着不能循环依赖。数仓管道天然是 DAG:
flowchart LR
A["ods_order 原始订单"] --> B["dwd_order 清洗明细"]
C["ods_user 原始用户"] --> D["dwd_user 清洗用户"]
B --> E["dws_user_order_1d 用户日汇总"]
D --> E
E --> F["ads_gmv_report GMV 报表"]
编排器读懂这张图后,调度规则很简单:一个任务的所有上游都成功了,它才变为「可调度」。dws_user_order_1d 必须等 dwd_order 和 dwd_user 都成功才会启动。
Airflow(Apache 项目)最大的特点是 DAG as code:用 Python 代码定义 DAG,而不是在界面上点。核心组件:
- Scheduler(调度器):不断解析所有 DAG,判断「哪些任务此刻该跑」(到点了 + 上游都成功了)。
- Executor(执行器):把该跑的任务真正发出去执行——可选
LocalExecutor(本机进程)、CeleryExecutor(分发到 worker 集群)、KubernetesExecutor(每个任务起一个 Pod)。 - Metadata DB:存所有任务的状态;Webserver:那张能看的 DAG 图和日志界面。
下面这段是一个每天凌晨跑的两步管道(Python,逐行注释,因为不假设你熟 Python 语法):
from airflow import DAG # DAG 定义入口
from airflow.operators.bash import BashOperator # 一种「执行 shell 命令」的算子;Operator=一类任务的模板
from datetime import datetime, timedelta # Python 标准库,处理时间
default_args = { # 本 DAG 所有任务的默认参数(一个字典/键值对)
"retries": 2, # 失败自动重试 2 次
"retry_delay": timedelta(minutes=5), # 每次重试间隔 5 分钟
}
with DAG( # 用 with 语法声明一个 DAG,块内定义的任务都归它
dag_id="dwd_order_daily", # DAG 的唯一名字
schedule="0 2 * * *", # cron 表达式:每天 02:00 触发(Airflow 2.4+ 用 schedule,旧版名为 schedule_interval)
start_date=datetime(2026, 1, 1), # 从哪天开始纳入调度,也是补数的起点
catchup=False, # 是否把 start_date 到现在的历史批次全部补跑;False=只跑最新一批
default_args=default_args, # 套用上面的默认参数
) as dag:
extract = BashOperator( # 任务 1:抽取 ODS 层订单
task_id="extract_ods_order",
# {{ ds }} 是 Airflow 模板变量,运行时渲染成「这一批次的数据日期」,如 2026-07-27
bash_command="spark-submit /jobs/extract_order.py --dt {{ ds }}",
)
transform = BashOperator( # 任务 2:清洗生成 DWD 层
task_id="transform_dwd_order",
bash_command="spark-submit /jobs/dwd_order.py --dt {{ ds }}",
)
extract >> transform # 声明依赖:extract 成功后才跑 transform(>> 是 Airflow 重载的「指向」运算符)
模拟运行结果(Web 界面上看到的一次调度):
DAG: dwd_order_daily run_id=scheduled__2026-07-27T02:00:00
extract_ods_order success (02:00:03 → 02:07:11)
transform_dwd_order success (02:07:12 → 02:15:44) # 上游成功后才启动
一句话解释:每天 02:00,Scheduler 发现 dwd_order_daily 到点,先让 extract_ods_order 变为可调度、被 Executor 领走执行;它成功后 transform_dwd_order 才启动;任一任务失败则按 retries=2 自动重试。注意命令里用的是 {{ ds }}(批次数据日期)而不是「当前系统时间」——这一点是下一节幂等与补数的关键。
DolphinScheduler(Apache 顶级项目,源自国内,原名 EasyScheduler)走另一条路:可视化拖拽画 DAG,不必写代码;去中心化的分布式架构,内置多种任务类型、补数、告警、租户隔离。
一句话对比(差异明显,用文字即可,不铺表格):Airflow 是「DAG as code」,工程师友好、版本可控、灵活但有门槛;DolphinScheduler 是「可视化」,数据分析师/运维友好、上手快,但复杂逻辑不如代码灵活。 选型看团队构成,不看谁「更先进」。
补数:把某段历史日期的任务重新跑一遍。为什么总要补数?① 上游数据修正了要重算;② 清洗逻辑改了要回刷历史;③ 新表首次上线要初始化历史分区。
补数能不能安全重跑,取决于任务是否幂等。
- 幂等(idempotency)定义:同一个任务用同样的输入重跑任意多次,结果都一致,不会产生重复或累加。
- 怎么做到幂等:按分区覆盖写——
INSERT OVERWRITE TABLE t PARTITION (dt='2026-07-27'),重跑就是把这个分区整块替换掉。 - 非幂等重跑的坑(真实事故高发区):
- 用
INSERT INTO(追加)而非OVERWRITE:补一次数据翻一倍,补两次翻三倍。 - 任务逻辑依赖
now()/ 系统当前时间,而不是批次数据日期(如 Airflow 的{{ ds }}):你想补 7 月 1 日的数据,任务却按「今天」去算,补出来全错。这就是 4.3 里强调「用{{ ds }}不用系统时间」的原因。
- 用
# Airflow 用命令行补 2026-07-01 到 07-07 这一周的数据
airflow dags backfill -s 2026-07-01 -e 2026-07-07 dwd_order_daily
# 每一天会以对应的 ds(07-01、07-02……)分别触发,因为任务按 ds 覆盖写分区,重跑安全
一句话解释:因为命令用 {{ ds }} 定位分区、且下游 INSERT OVERWRITE 覆盖写,所以补哪天就精确重刷哪天的分区,跑几次结果都一样——幂等是补数敢按按钮的前提。
概念都齐了,回到第一节的那句话做个逐维度对照(这两个概念高频混淆,值得用表格拆开——差异明显的地方本可用文字,但这里正因为都叫「调度」才需要并排看):
| 维度 | 资源调度(YARN / K8s) | 工作流编排(Airflow / DolphinScheduler) |
|---|---|---|
| 回答什么问题 | 给多少 CPU/内存、在哪台机器上跑 | 什么时候触发、依赖谁先跑完 |
| 核心抽象 | Container / Pod(一份算力额度) | DAG(任务的依赖图) |
| 调度粒度 | 一个作业内部的进程(Driver/Executor) | 一个作业作为一个整体的节点 |
| 时间尺度 | 秒级分配、随作业存续 | 分钟级触发、按天/小时周期 |
| 失败语义 | 容器挂了重新申请一个 | 任务失败按 retries 重试、告警、补数 |
| 典型产品 | YARN、K8s(+ Volcano 项目) | Airflow、DolphinScheduler |
它们串在一条链上:编排器(第 4 节)按定时和依赖决定「现在该跑 dwd_order 了」→ 发起 spark-submit → 资源调度器(第 2/3 节)给这个 Spark 作业分配容器 → Driver/Executor 真正开算(05 篇)→ 结果落表 → 编排器据此触发下游。一个定顺序,一个给算力,各管一段,缺一不可。
边界
资源调度侧:
- 资源碎片:每台机器剩的零头都不够放下一个 Container/Pod,总量够但拼不出一整块 → 大作业申请不到、利用率虚低。缓解靠更紧的装箱(bin-packing)和更合理的容器规格,别把 Executor 配得又大又不整。
- 抢占(preemption):高优先级队列可以抢占低优先级正在跑的 Container。代价是被抢的任务丢失已算的进度、要重来——抢占提升了高优作业的响应,却是拿低优作业的重算成本换的。
- gang scheduling 缺失:K8s 原生调度器跑批会资源死锁,生产跑 Spark on K8s 记得上 Volcano 项目 / YuniKorn。
工作流编排侧:
- DAG 循环依赖:建模时不小心让「表 A 依赖 B、B 又依赖 A」,编排器要么报错、要么两个任务永远
pending。根因在数仓建模,不在工具——分层清晰(ODS→DWD→DWS→ADS 单向流)就不会成环。 - 补数踩坑:非幂等重跑(见 4.5)是头号杀手;其次是补数并发过大(一次性把一个月每天的任务全放出去,直接压垮集群),要限并发、按
dt有序回填。 - 调度延迟:Airflow Scheduler 有解析/心跳间隔,任务不是「准点到秒」触发,分钟级延迟是常态;资源不够时还会排队 pending。别把编排器当实时系统。
何时不该用:
- 别用工作流编排器做秒级/高频实时触发——它是面向批量、分钟级周期的。实时链路请用流处理(见 07 篇),二者定位不同。
- 别把「绝对不能被打断」的关键作业放进会被抢占的低优队列——给它独立队列或资源池。
- 小规模、单机就能跑的作业别硬上 YARN/K8s 全家桶——分布式资源调度的运维成本很可能大于收益,一台大内存机器 + 本地 Spark/
cron就够了。
资源调度侧:
- EMR + YARN(官方叫「存算一体」):火山引擎 E-MapReduce(EMR) 是托管的 Hadoop 生态集群,官方把「HDFS 存储 + YARN 计算同集群」称为存算一体方案——即 EMR 上的资源调度正是本篇第 2 节的 YARN(RM/NM/AM/Container)。官方把 Spark 列为支持引擎之一,你在 EMR 上跑 Spark 时,2.5 那张时序图就是幕后发生的事。(注:官方文档未使用「Spark on YARN」这个字面词,那是社区惯用叫法。来源 docs/6491/149821。)
- 存算分离 + Proton:把底层存储从 HDFS 换成对象存储 TOS,官方称为存算分离方案,并为此提供自研加速引擎 Proton——对应本篇第 3 节「数据在对象存储、计算随起随销」,也是你用过的云上 Iceberg/LAS 湖仓弹性的底座。(来源同上 docs/6491/149821。)
- EMR on VKE(K8s 化,官方产品名):这正是第 3 节说的 K8s 化调度落地——官方形态叫 EMR on VKE,明确支持「在线业务系统和离线大数据共享一套 VKE(K8s)集群」,靠弹性伸缩、动态超分提升利用率,其上可跑 Spark 和 Ray。第 3 节的「混部 + 弹性」在这里就是产品。(来源 docs/6491/1382645。)
- EMR Serverless:更进一步的 Serverless 形态(官方产品名是 EMR Serverless,不叫「Serverless Spark」),主打「弹性扩展、秒级资源响应、按量付费」;作业通过统一的作业提交接口递交,背后由 Kyuubi Server 承接作业提交与调度。(来源 docs/6491/1382644。)
工作流编排侧:
- DataLeap 的「数据开发」:火山引擎侧承担编排的是大数据研发治理套件 DataLeap 里的数据开发(DataStudio)模块。本篇的能力它都有,但官方术语和开源栈不同:
- 依赖编排:官方叫工作流(一个工作流里拖拽多个子节点统一编排),依赖类型分「上游依赖 / 同频依赖 / 跨周期自依赖」——它是 DAG 式编排,但官方文档不用「DAG / 有向无环图」这个词。
- 定时:官方叫周期调度 / 周期实例,支持准实时到月级。
- 失败重试:官方参数叫失败重跑次数 / 重跑时间间隔。
- 补数:官方功能叫数据回溯,正文用补数据 / 重跑描述(注意:官方不叫「数据补录」,别记混)。 (来源 docs/6260/65374、docs/6260/75005。)
- EMR 内置 DolphinScheduler:传统 EMR on ECS 集群里内置了本篇 4.4 提到的开源 DolphinScheduler 做作业调度——你在 EMR 上见到的那个调度组件不是别的,就是它。
云数仓 ByteHouse 侧:
- ByteHouse 定位是查询与存储引擎(对应 03/04/10 篇),没有通用的跨任务 DAG 编排——那是 DataLeap 的活。但它自带一个小而具体的编排能力值得一提:异步物化视图的刷新调度,用
REFRESH ASYNC EVERY (INTERVAL ...)声明定时刷新,后台由DaemonManager起线程执行。- 坑(官方写明,正好呼应第 6 节的「失败语义」):这个刷新间隔不严格保证(一次刷新耗时超过间隔就顺延到下个周期),而且没有重试机制——失败只能等下一个周期。这是一个「编排能力很轻、失败语义很弱」的反例,别拿它当 DataLeap 那种正经编排器用。(来源 docs/6517/1298237。)
落到你的日常:当你在平台上「配一个每天跑的 SQL、连上上游、设失败重跑次数」,你操作的是 DataLeap 数据开发 = 工作流编排层;当任务真正跑起来、后台去 EMR / EMR on VKE 申请 CPU/内存把引擎拉起来,那是 YARN / VKE = 资源调度层 在干活。两层,正是本篇讲的两件事。