跳到主要内容

实时任务

最近更新 2026/10/03

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: 生产环境值}

参数校验规则(仅在调试/执行时校验,保存和发布时不校验):

  1. 任务参数:若在 SQL 中声明了 ${abc} 但未定义值,调试/执行时会拦截并要求用户赋值
  1. 空间参数:若 SQL 中引用的空间参数被删除,调试/执行时会拦截并提示参数不存在

2. 调试实时任务​

2.1 调试功能说明​

调试功能用于在 开发环境 中验证实时任务的正确性,无需发布到生产环境即可测试任务逻辑。

点击「调试」按钮时,可选择:

  • 仅保存:仅触发保存校验流程
  • 保存并调试:触发保存校验 + 调试校验流程

调试不是发布前的必经流程。用户保存后可直接发布,无需先调试。

2.2 校验流程​

保存校验:

  1. SQL 语法校验

  2. 表权限校验(来源表和目标表的读写权限)

  3. 参数解析与记录

调试/执行校验(在保存校验基础上增加):

  1. 参数完整性校验 — 所有 ${XX} 参数必须有值
  2. 空间参数存在性校验
  3. 执行通道状态校验

2.3 调试执行​

调试成功后,任务会在开发环境的执行通道中运行。可以:

  • 跳转Flink UI 查看任务运行状态
  • 查看实时日志
  • 查看运行指标
  • 手动停止调试任务

3. 保存与发布实时任务​

3.1 保存​

实时任务的发布(上线)是单向的 — "已发布"的实时任务不能回退到"未发布"状态(与离线任务流可双向切换不同)。

任务"待发布"时保存:

项目说明
存在环境仅保存在「开发环境」
上线状态"待发布"
版本开发环境展示最新的版本;生产环境无
可操作编辑、发布上线、删除

任务"已发布"时保存:

项目说明
存在环境新版本内容仅保存在「开发环境」,需重新发布到生产
上线状态"已发布"(有变更)
版本开发环境展示最新的;生产环境仍为上一次发布的版本
可操作编辑、重新发布、运行、删除

3.2 保存并发布​

发布操作会将开发环境的实时任务版本同步到生产环境。

首次发布:

  • 开发环境和生产环境同时更新版本

  • 发布成功后状态变为"已发布",需点击「启动」运行首个实例

  • 发布失败则保持"待发布"状态

已发布任务重新发布:

  • 开发环境和生产环境版本均更新

  • 状态保持"已发布"

  • 发布成功后点击「启动」使用新版本运行实例

3.3 发布校验​

发布时会进行以下校验:

  1. 未保存的任务不可发布(发布按钮置灰)
  2. 生产环境表权限校验
  3. 参数值完整性校验
  4. 表结构一致性校验

4. 运行(启动)实时任务​

4.1 生产环境详情页​

任务发布到生产环境后,可在 生产环境详情页 查看和管理任务。详情页包含以下模块:

模块说明
实时任务基本信息任务名称、负责人、引擎版本、变更状态、创建/修改时间、描述
运行配置执行通道选择(资源配置)、运行配置参数
实例内容最新实例的 SQL 内容、运行配置、表血缘、任务参数值
运行记录实例完整生命周期的事件消息
日志包含运行日志和异常日志信息
快照Checkpoint/Savepoint 的总数、成功/失败统计及历史详情

4.2 运行参数配置​

启动 实时任务前需要完成以下配置:

4.2.1 参数赋值​

对于「代码解析参数」,需要在启动时赋值,参数值记录在实例级别。

4.2.2 运行配置​

运行资源配置

配置项说明

通道(资源池)选择

Channel 名称:系统会根据任务 Flink 引擎版本自动过滤可用的通道

只能选择 "已启用" 的通道

实时任务执行只能选择 运行模式 = Application 的通道(资源池)

计算模式
  • 资源模式分为:计算型 / 通用型。
  • 计算型(1:2):每核配2G内存,CPU性能更强,适合密集计算;通用型(1:4):每核配4G内存,CPU与内存均衡,适合多数业务
  • 可选的资源模式由所选资源池绑定的可选规格决定
并行度

资源规格在通道级已定义,实例提交执行时选择 并行度 数量。

并行度是 Flink 对任务并行切分的描述。Slot 数目是对单个 TaskManager 的资源切分粒度。

运行参数配置

类别参数说明默认值
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重启的策略,则会按照配置的重启策略进行重启。

该参数取值如下:

  • Failure Rate(基于失败率重启):设定检测时间间隔、最大失败次数、每次重启间隔
  • Fixed Delay(固定间隔重启):设定尝试重启次数、每次重启间隔
  • No Restarts(不重启):作业失败后不自动重启
Fixed Delay

4.2.3 启动策略​

用户启动/重启任务时可选择是否从已保存的快照状态启动:

启动策略说明
有状态启动从已存在的有效状态(Checkpoint / Savepoint)恢复。可选择从最新状态恢复,或从历史快照列表中选择
无状态启动不包含初始状态,全新启动

注意: 任务不支持按历史版本启动。任意 Savepoint 都可选择,但 Schema 兼容性需用户自行负责,若 Schema 变更可能导致运行时报错。

4.2.4 提交启动​

在启动前会执行以下校验,校验成功后启动实时任务:

  1. 任务存在性校验

  2. 内容合法性校验

  3. 版本一致性校验

  4. 通道存在性校验

  5. 快照存在性校验(有状态启动时)

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 快照管理​

字段说明
快照 IDStreampark 返回的 Checkpoint ID
状态快照的当前状态
快照类型

Checkpoint:Flink 自动创建的"自动存档",用于故障恢复,轻量、快速、自动化

Savepoint:用户手动触发的"手动存档点",用于有计划地停止与重启(代码更新、集群扩缩容、版本升级等)

创建方式手动(用户发起,仅 Savepoint)/ 自动(根据 Checkpoint 配置自动生成)
触发时间快照触发的时间
持续时间快照创建耗时
来源实例快照所属的实例 ID

7. 管理页元数据​

在实时任务管理列表页,可查看以下元数据信息:

字段说明
实时任务名称任务唯一标识
备注任务描述信息
上线状态已发布 / 待发布
运行状态当前版本最新实例的运行状态
引擎版本如 Flink 1.16
负责人任务负责人
上线时间开发环境发布上线生产的时间
启动时间最新实例的启动时间
最后修改时间最近一次修改的时间
最后修改人最近一次修改的操作人
操作编辑、状态管理操作、任务详情跳转
这篇文档对你有帮助吗?