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

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

3. 查询已写入的数据
作业跑完就能直接查。这是一条不涉及复杂计算的简单查询,走共享的 Trino 引擎执行,不需要为本空间单独申请 Spark 查询引擎——省去一次配额申请,直接拿到结果:
sql
SELECT * FROM iceberg_lake.site_demo_ws10.daily_sim_orders ORDER BY dt
想直接查数而不是写调度作业,看数据查询;想了解作业跑挂了怎么办,看AI 智能诊断。
看全部作业类型 →