实时任务
AE 数据开发平台通过深度集成 Apache Flink 强大的流式计算能力,提供性能更强的流式数据同步方案,并支持对流数据进行实时分析与探索,为企业打造面向未来的 流批一体、高实时性、低运维成本 的企业级数据开发平台。
1. 创建实时任务
1.1 新建入口
在数据开发平台的「开发」模块中,点击「实时任务」页签右上角的 「+ 实时任务」按钮,进入创建页面。
1.2 填写基本信息
实时任务支持双环境(开发环境 / 生产环境),基本信息为两个环境共用。
| 字段 | 说明 | 备注 |
|---|---|---|
| 名称 | 空间内唯一,不超过 48 字符,支持中文、英文、数字、下划线 | 必填 |
| Flink 版本 | 实时任务绑定的 Flink 引擎版本,从下拉列表选择 | 创建后不可修改 |
| 负责人 | 实时任务的负责人,创建人默认为负责人 | 必填 |
| 备注 | 任务描述信息,不超过 200 字符 | 选填 |
1.3 开发环境编辑页面
创建完成后进入 开发环境 编辑页,主要包含以下区域:
1.3.1 代码编辑区
-
支持 Flink SQL 语法高亮与自动补全
-
支持参数化配置,使用
${参数名}格式引用参数 -
编辑器会自动解析代码中的参数,并在下方「使用的参数」区域展示
-
编辑器同时还会解析表的血缘关系,展示 Flink SQL 中涉及的表的来源、去向,connector 类型等信息
1.3.2 任务参数管理(侧边栏)
实时任务支持以下两种参数类型:
| 参数类型 | 说明 |
|---|---|
| 空间参数 | 支持 预置空间参数 和 自定义空间参数 |
| 任务参数 | 实时任务级别 隔离,支持「纯文本」和「表达式」两种类型 |
参数新增方式:
-
方式一:手动创建 — 在侧边栏中手动添加参数,提交任务时会将参数值解析并传递给执行引擎
-
方式二:代码解析 — 在 Flink SQL 中直接使用
${abc}格式定义参数,系统会自动解析并在侧边栏展示
动态参数:
当需要在不同环境使用不同参数值时,可以使用 动态参数,在运行时指定当前环境取的对应值。
${env.dev: 开发环境值, env.product: 生产环境值}
参数校验规则(仅在调试/执行时校验,保存和发布时不校验):
- 任务参数:若在 SQL 中声明了
${abc}但未定义值,调试/执行时会拦截并要求用户赋值
- 空间参数:若 SQL 中引用的空间参数被删除,调试/执行时会拦截并提示参数不存在
2. 调试实时任务
2.1 调试功能说明
调试功能用于在 开发环境 中验证实时任务的正确性,无需发布到生产环境即可测试任务逻辑。
点击「调试」按钮时,可选择:
- 仅保存:仅触发保存校验流程
- 保存并调试:触发保存校验 + 调试校验流程
调试不是发布前的必经流程。用户保存后可直接发布,无需先调试。
2.2 校验流程
保存校验:
-
SQL 语法校验
-
表权限校验(来源表和目标表的读写权限)
-
参数解析与记录
调试/执行校验(在保存校验基础上增加):
- 参数完整性校验 — 所有
${XX}参数必须有值 - 空间参数存在性校验
- 执行通道状态校验
2.3 调试执行
调试成功后,任务会在开发环境的执行通道中运行。可以:
- 跳转Flink UI 查看任务运行状态
- 查看实时日志
- 查看运行指标
- 手动停止调试任务
3. 保存与发布实时任务
3.1 保存
实时任务的发布(上线)是单向的 — "已发布"的实时任务不能回退到"未发布"状态(与离线任务流可双向切换不同)。
任务"待发布"时保存:
| 项目 | 说明 |
|---|---|
| 存在环境 | 仅保存在「开发环境」 |
| 上线状态 | "待发布" |
| 版本 | 开发环境展示最新的版本;生产环境无 |
| 可操作 | 编辑、发布上线、删除 |
任务"已发布"时保存:
| 项目 | 说明 |
|---|---|
| 存在环境 | 新版本内容仅保存在「开发环境」,需重新发布到生产 |
| 上线状态 | "已发布"(有变更) |
| 版本 | 开发环境展示最新的;生产环境仍为上一次发布的版本 |
| 可操作 | 编辑、重新发布、运行、删除 |
3.2 保存并发布
发布操作会将开发环境的实时任务版本同步到生产环境。
首次发布:
-
开发环境和生产环境同时更新版本
-
发布成功后状态变为"已发布",需点击「启动」运行首个实例
-
发布失败则保持"待发布"状态
已发布任务重新发布:
-
开发环境和生产环境版本均更新
-
状态保持"已发布"
-
发布成功后点击「启动」使用新版本运行实例
3.3 发布校验
发布时会进行以下校验:
- 未保存的任务不可发布(发布按钮置灰)
- 生产环境表权限校验
- 参数值完整性校验
- 表结构一致性校验
4. 运行(启动)实时任务
4.1 生产环境详情页
任务发布到生产环境后,可在 生产环境详情页 查看和管理任务。详情页包含以下模块:
| 模块 | 说明 |
|---|---|
| 实时任务基本信息 | 任务名称、负责人、引擎版本、变更状态、创建/修改时间、描述 |
| 运行配置 | 执行通道选择(资源配置)、运行配置参数 |
| 实例内容 | 最新实例的 SQL 内容、运行配置、表血缘、任务参数值 |
| 运行记录 | 实例完整生命周期的事件消息 |
| 日志 | 包含运行日志和异常日志信息 |
| 快照 | Checkpoint/Savepoint 的总数、成功/失败统计及历史详情 |
4.2 运行参数配置
启动 实时任务前需要完成以下配置:
4.2.1 参数赋值
对于「代码解析参数」,需要在启动时赋值,参数值记录在实例级别。
4.2.2 运行配置
运行资源配置
| 配置项 | 说明 |
|---|---|
通道(资源池)选择 | Channel 名称:系统会根据任务 Flink 引擎版本自动过滤可用的通道
|
| 计算模式 |
|
| 并行度 | 资源规格在通道级已定义,实例提交执行时选择 并行度 数量。
|
运行参数配置
| 类别 | 参数 | 说明 | 默认值 |
|---|---|---|---|
| Checkpoint 设置 | execution.checkpointing.interval | 检查点间隔 | 180 秒 |
execution.checkpointing.timeout | 检查点超时时间 | 3 分钟 | |
execution.checkpointing.min-pause | 检查点间最短间隔 | 60 秒 | |
状态数据过期 | table.exec.state.ttl | 状态数据过期时间,0 表示永不过期 当数据首次进入系统并被处理后,它会存储在状态内存中。当下一次相同主键的数据到来时,系统会使用之前存储的状态数据进行计算,并更新其访问时间。这一过程是实时计算的核心,因为它依赖于数据的持续流动。如果数据在设定的TTL时间窗口内未被再次访问,它将被系统视为过期,并从状态存储中清除。通过合理设置TTL的值,不仅可以维持计算的精确性,还能及时清理陈旧数据,有效减少状态内存的占用,进而降低系统内存负担,提升计算效率和系统稳定性。 | 0 |
| 重启策略 | restart-strategy.type | 重启策略类型 只有在没配置重启策略的情况下,Flink才会根据系统检查点开启与否来决定是否要重启作业(如果系统检查点开启则会按Fixed Delay设置的固定间隔重启作业,未开启则不重启作业)。如果配置了Flink重启的策略,则会按照配置的重启策略进行重启。 该参数取值如下:
| Fixed Delay |
4.2.3 启动策略
用户启动/重启任务时可选择是否从已保存的快照状态启动:
| 启动策略 | 说明 |
|---|---|
| 有状态启动 | 从已存在的有效状态(Checkpoint / Savepoint)恢复。可选择从最新状态恢复,或从历史快照列表中选择 |
| 无状态启动 | 不包含初始状态,全新启动 |
注意: 任务不支持按历史版本启动。任意 Savepoint 都可选择,但 Schema 兼容性需用户自行负责,若 Schema 变更可能导致运行时报错。
4.2.4 提交启动
在启动前会执行以下校验,校验成功后启动实时任务:
-
任务存在性校验
-
内容合法性校验
-
版本一致性校验
-
通道存在性校验
-
快照存在性校验(有状态启动时)
5. 实时任务运行状态管理
5.1 任务实例管理状态
| 状态 | 可流转到的状态 | 可执行操作 | 说明 |
|---|---|---|---|
| 未发布 | — | 编辑、发布上线、删除 | 实时任务仅存在于开发环境 |
| 未启动 (Not Started) | — | 编辑、发布上线、删除 | 已发布生产环境,但尚未运行 |
| 启动中 (Starting) | 运行中、失败、停止中、杀死中 | 编辑、发布上线、停止、Kill | 正在提交任务到执行引擎 |
| 运行中 (Running) | 已停止、停止中、杀死中、失败 | 编辑、发布上线、停止、Kill | 任务正常运行 |
| 失败 (Failed) | 启动中(新实例) | 编辑、发布上线 | 实例终态,可重启新实例 |
| 停止中 (Stopping) | 已停止、杀死中、失败 | 编辑、发布上线、Kill | 正在执行停止操作 |
| 已停止 (Stopped) | 启动中(新实例) | 编辑、发布上线 | 实例终态,可重启新实例 |
| 杀死中 (Killing) | 已停止、失败 | 编辑、发布上线 | 正在强制终止任务 |
注意: Stopping 和 Killing 超时会发送告警信息。Stopping 过程中支持二次停止操作。
5.2 Stop 与 Kill 的区别
| 操作 | 说明 |
|---|---|
| 停止 (Stop) | 优雅停止,可选择是否触发 Savepoint 保存状态。停止后可从 Savepoint 恢复 |
| Kill | 强制终止,不保存状态。适用于 Stop 无响应或任务卡死的场景 |
- 停止任务时: 自动触发 Savepoint(存储路径无需用户定义)
- 启动/重启任务时: 可选择是否从 Checkpoint/Savepoint 恢复
6. 生产环境详情页
6.1 实例内容
最新实例信息,包含
- 脚本 SQL 内容
- 实例的运行配置
- 实例的表血缘
- 实例包含的任务参数值
6.2 运行记录
运行记录按照 Streampark Job Instance ID 记录每个实例的完整生命周期事件:
- 默认展示最新实例的事件消息,按实例创建时间倒序排列
- 用户可展开查看所有实例的事件消息
- 支持时间轴控制全局信息展示
6.3 日志
| 分类 | 内容 | 说明 | |
|---|---|---|---|
运行日志 | JobManager 日志 | 任务调度、Checkpoint、状态管理 | 单个 JM 实例的日志 |
| TaskManager 日志 | 数据处理、网络通信 | 多个 TM 实例的日志聚合 | |
| 异常信息 | ERROR/WARN 日志聚合 | 快速定位问题的入口 |
6.4 快照管理
| 字段 | 说明 |
|---|---|
| 快照 ID | Streampark 返回的 Checkpoint ID |
| 状态 | 快照的当前状态 |
| 快照类型 | Checkpoint:Flink 自动创建的"自动存档",用于故障恢复,轻量、快速、自动化 Savepoint:用户手动触发的"手动存档点",用于有计划地停止与重启(代码更新、集群扩缩容、版本升级等) |
| 创建方式 | 手动(用户发起,仅 Savepoint)/ 自动(根据 Checkpoint 配置自动生成) |
| 触发时间 | 快照触发的时间 |
| 持续时间 | 快照创建耗时 |
| 来源实例 | 快照所属的实例 ID |
7. 管理页元数据
在实时任务管理列表页,可查看以下元数据信息:
| 字段 | 说明 |
|---|---|
| 实时任务名称 | 任务唯一标识 |
| 备注 | 任务描述信息 |
| 上线状态 | 已发布 / 待发布 |
| 运行状态 | 当前版本最新实例的运行状态 |
| 引擎版本 | 如 Flink 1.16 |
| 负责人 | 任务负责人 |
| 上线时间 | 开发环境发布上线生产的时间 |
| 启动时间 | 最新实例的启动时间 |
| 最后修改时间 | 最近一次修改的时间 |
| 最后修改人 | 最近一次修改的操作人 |
| 操作 | 编辑、状态管理操作、任务详情跳转 |

