主题
Flink SQL 作业
场景 跟前面几种批作业不同,Flink SQL 是 7×24 常驻的:启动之后就一直跑,不会自己"结束",直到手动停止。典型如实时聚合看板、把消息/变更流实时落进湖仓。
能干什么
| 你想做的事 | 配置成 Flink SQL 长这样(示意) |
|---|---|
| 现场造数据验证环境 | CREATE TABLE src (...) WITH ('connector'='datagen', ...) |
| 消费 Kafka 写入湖仓 | CREATE TABLE src (...) WITH ('connector'='kafka-ds', 'ds.ref'='kafka_xxx', ...) |
| MySQL 变更实时入湖 | CREATE TABLE src (...) WITH ('connector'='mysql-cdc-ds', 'ds.ref'='mysql_xxx', ...) |
| 写入 Iceberg 湖表 | INSERT INTO iceberg_lake.库.表 SELECT ... FROM src |
亮点
- 建表 + INSERT 就是一个流作业——纯 Flink SQL,不用写 Java/Scala,跟批作业里写 SQL 的体验是一致的。
- 连接器开箱即用,凭证不进 SQL 明文——Kafka、MySQL CDC、Iceberg 连接器都已经内置;SQL 里只写
'ds.ref'='<数据源名>'引用已注册的数据源,真实的地址/账号/密码由平台在运行时注入,不用也不能把凭证写进作业 SQL 里。 - 启动/停止是独立的生命周期,不是"运行一次"——作业详情页有专属的「实时集成部署」面板:状态是运行中/已停止,不是批作业那种成功/失败的运行历史列表;按「启动」部署,按「停止(savepoint)」下线。
- 停止即做 savepoint,下次自动续跑——点「停止」会先触发一次 savepoint 再优雅停机;下次「启动」自动从这个 savepoint 恢复位点,不用从头重放数据(exactly-once 续传)。
- 逐表鉴权,连自己的流作业也不例外——往一张从未在本空间被正式创建过的表
INSERT INTO,会被引擎侧鉴权直接拒绝、作业启动失败,哪怕执行者是空间管理员;稳妥做法是先用别的方式(SQL 作业/查询)把目标表建好,流作业只管写,不替你悄悄用信任更高的身份建表。 - 资源跟批作业共用一个池子,但流作业不排队——底层也是本空间的 Flink 集群,占的是同一份「批作业资源」配额;区别是配额不够时批作业会排队等,流作业直接拒绝启动。
完整操作示例:用 datagen 现场生成数据,持续写入 Iceberg
下面是一次真实的流水线:不依赖任何外部系统(不用先搭 Kafka),直接用 Flink 内置的 datagen 连接器现场生成模拟传感器数据,持续写入 Iceberg,跑一段时间后查询验证数据真的在不断进来,最后演示停止(savepoint)。
1. 先建好目标表
Flink 作业如果直接 INSERT INTO 一张从未创建过的表,会在鉴权层被拒绝——鉴权检查的是"这张表有没有被本空间正式创建过",不是语法对不对。所以先把目标表建好:
sql
CREATE TABLE IF NOT EXISTS iceberg_lake.site_demo_ws10.flink_datagen_sensor (
device_id STRING,
temperature DOUBLE,
ts TIMESTAMP
) USING iceberg2. 配置 Flink SQL 作业:datagen 生成 + 写入 Iceberg
新建「Flink SQL」作业,选上「Iceberg catalog」(选完 SQL 里就能直接用 catalog.库.表 引用湖表),SQL 分两段:先建一张 datagen 源表(限定字段范围,生成的数据更好读),再 INSERT INTO 写进第一步建好的湖表。
sql
CREATE TABLE src (
device_id STRING,
temperature DOUBLE,
ts AS PROCTIME()
) WITH (
'connector' = 'datagen',
'rows-per-second' = '2',
'fields.device_id.length' = '6',
'fields.temperature.min' = '18.0',
'fields.temperature.max' = '30.0'
);
INSERT INTO iceberg_lake.site_demo_ws10.flink_datagen_sensor
SELECT device_id, temperature, ts FROM src;
3. 启动,查看运行状态
点「启动」部署,平台会拉起一个专属这个作业的 Flink 集群。运行起来后,面板上能看到部署状态、Flink 原生状态、运行时长,以及 checkpoint 完成/失败次数、最近一次 checkpoint 的时间和大小——这些都是持续更新的,不是跑完就定格的批作业日志。

4. 查询验证:数据在持续流入
跑了几分钟后去查一下,first_ts 到 last_ts 跨度和当前累计行数,能直接说明数据是随时间持续写入的,不是一次性灌进去的:
sql
SELECT COUNT(*) AS cnt, MIN(ts) AS first_ts, MAX(ts) AS last_ts
FROM iceberg_lake.site_demo_ws10.flink_datagen_sensor
5. 停止(savepoint)
不需要了就点「停止(savepoint)」,平台会先触发一次 savepoint 再停掉集群。状态变回「已停止」,Flink JobID 保留可查;下次「启动」会自动从这次的 savepoint 续跑,不丢位点。

想批量同步外部数据库(非实时),看离线集成(DATA_SYNC)作业;想了解作业跑挂了怎么办,看AI 智能诊断。
看全部作业类型 →