全文共 5,179 字 预计阅读 15 分钟
bg

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-b 40%),队列之间硬隔离,队列内部再按 FIFO/优先级排。适合多租户按部门分账这是 Apache Hadoop 当前版本的默认调度器——yarn-default.xmlyarn.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 脚本串起来。痛点很快暴露:

  1. 依赖靠猜:下游任务不知道上游好没好,只能「估上游 3 点跑完,我排 4 点」,靠 sleep 和错峰硬扛——上游一延迟,下游就读到半成品数据。
  2. 失败没兜底:脚本挂了没人知道,没有自动重试、没有告警。
  3. 补历史数据是灾难:要重算上个月的数据,只能手改脚本里的日期,跑一天改一次。
  4. 依赖关系没人看得清:几百个任务的先后关系散在几十个 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_orderdwd_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 = 资源调度层 在干活。两层,正是本篇讲的两件事。

Back to Blog