Skip to content

Spark SQL 作业

场景 数据仓库里的离线批处理——一次性跑完就结束,不是流作业那种 7×24 常驻的。典型如每天凌晨算一版报表、把明细表汇总成宽表、建一张新表并灌入数据。

能干什么

你想做的事写成 Spark SQL 长这样(示意)
汇总昨天的订单INSERT OVERWRITE dw.order_daily SELECT dt, sum(amt) FROM ods.orders WHERE dt='2026-07-26' GROUP BY dt
建一张新宽表CREATE TABLE dw.user_wide (...) USING iceberg
建表并灌数据CREATE TABLE dw.t USING iceberg AS SELECT ...(CTAS,或先建表再 INSERT)
跨库关联SELECT ... FROM lake.a JOIN olap.b ON ...

亮点

  • 跨源关联——一条 SQL 里可以引用湖仓(Iceberg)、OLAP(Doris)、外部数据库(MySQL / PostgreSQL)等区域内任意已注册的数据源,不用先把数据倒来倒去。
  • 逐表逐列鉴权——SQL 每引用到一张表都会先做权限校验,无权限直接拒绝;隐私列自动脱敏。作业能拿到足够权限干活,但越不了权。
  • 跑的时候能看——运行中直接看 Spark 原生 UI,DAG、每个 stage/task 的进度和指标一目了然,不是提交完就两眼一抹黑。
  • 跑完还能查——每个空间按需生成独立的历史 UI,事后也能回看 DAG 和报错栈,方便定位问题。
  • 和调度打通——可以和其它类型的作业一起编排进 DAG 工作流,按上游触发规则联动;跑失败了旁边就是「AI 智能诊断」和一键重跑。
  • 资源互不挤占——每个空间独立配额与资源池,不会被别的空间的作业占满。

完整操作示例:每天定时写入 + 查询验证

下面是一个连贯的真实作业:每天定时写入一条模拟订单,配好告警,再去查验证数据真的进来了。

前置条件 新空间默认没有批作业计算配额,批作业按配额校验运行。先在「空间管理 → 资源管理」为批/流作业申请 CPU、内存和并发数,提交给平台审批,通过后才能跑批作业。

申请批作业执行资源

1. 建表 + 每次写入一条模拟数据

作业首次跑会建一张按 dt 分区的表,之后每次调度都往里插一条模拟订单(dt 取当天日期),这样才是"每天定时写入"的真实场景,而不是一次性灌几行静态数据。按天分区之后,以后按日期过滤查询能直接跳过不相关分区,也方便按天做生命周期管理:

sql
CREATE TABLE IF NOT EXISTS iceberg_lake.site_demo_ws10.daily_sim_orders (
  order_id BIGINT,
  amount DECIMAL(10,2),
  dt DATE
) USING iceberg
PARTITIONED BY (dt);

INSERT INTO iceberg_lake.site_demo_ws10.daily_sim_orders
SELECT CAST(900 + FLOOR(rand() * 100) AS BIGINT) AS order_id,
       CAST(ROUND(50 + rand() * 200, 2) AS DECIMAL(10,2)) AS amount,
       current_date() AS dt

2. 配上周期调度和告警

调度方式选「周期调度」,填 Cron 0 0 2 * * ?(每天凌晨 2 点跑一次);同时打开「告警设置」,失败或超时都会通知到指定成员,不用等有人发现作业早就没在跑了才后知后觉。

作业的周期调度与告警配置

配好之后,作业详情页能直接看到调度方式、Cron 表达式,以及历次运行记录:

周期调度作业与运行历史

3. 查询已写入的数据

作业跑完就能直接查。这是一条不涉及复杂计算的简单查询,走共享的 Trino 引擎执行,不需要为本空间单独申请 Spark 查询引擎——省去一次配额申请,直接拿到结果:

sql
SELECT * FROM iceberg_lake.site_demo_ws10.daily_sim_orders ORDER BY dt

查询到当天写入的模拟数据

想直接查数而不是写调度作业,看数据查询;想了解作业跑挂了怎么办,看AI 智能诊断

看全部作业类型 →