在复杂的生产环境中,保障 DataWorks 任务的稳定,不仅需要关注任务自身的调度状态,更需要深入洞察任务底层依赖的引擎(如 Flink、StarRocks)的实时健康状况。DataWorks AI 助理(DataClaw)将平台的观测能力从"任务调度层"深化到"引擎依赖层",支持对 Flink 作业进行主动、周期性的健康度巡检,并将巡检结果推送至钉钉等 IM 端。当 Flink 任务存在数据延迟、链路异常、资源水位过高等隐性问题时,只需配置巡检任务,AI 助理即可自动完成指标采集、异常判定与结果推送,实现任务全生命周期信息的统一收口,异常解读与后续处理。
本文以 Flink 流作业搭建的实时数仓为例,介绍如何通过 AI 助理了解当前实时数仓存在的问题。
前提条件
使用 AI 助理巡检 Flink 作业前,需满足以下条件:
- 已创建并启用 DataWorks AI 助理实例。如未创建,请参见使用 AI 助理服务。
- 已启用阿里云实时计算 Flink 专家套件。启用套件参见:启用专家套件并确认监控能力。
- 确保DataWorksAI助理执行身份具备对相应监控资源对象的访问权限。
- 若需在报警中 @ 责任人,需配置阿里云账号与 IM 端账号映射。详情参见:配置账号映射。
巡检能力
DataWorks 上 Flink 任务的执行稳定性,除了任务本身的调度运行状态外,还需要关注引擎执行层的隐性问题,如数据延迟、链路异常、资源水位过高等。为了及时发现这些问题,Flink 提供了两类告警机制:根据指标数值变化进行指标告警;根据事件是否发生进行事件告警。
告警类型 |
说明 |
示例 |
指标告警 |
根据指标数值变化触发 |
CPU 使用率、延迟、数据量等 |
事件告警 |
根据事件是否发生触发 |
作业失败、ECS 宕机等 |
为全面了解任务执行情况,除了 Flink 文档提到的云监控或 ARMS 监控方式外,您也可以通过 DataWorks AI 助理定时巡检任务,将这些告警对接至 DataWorks,通过钉钉等 IM 端接收报警,对任务全生命周期相关信息统一收口至DataWorks,并通过DataWorks AI 助理进一步解读异常。
- 可检查的 Flink 指标(27项):参见 Flink 监控指标。
您也可以通过与AI助理对话实时获取支持的指标,开启阿里云实时计算 Flink 专家套件后,可通过对话了解各指标意义。
- 可检查的 Flink 系统事件(2类):
事件 |
级别 |
说明 |
作业运行失败(JOB_FAILED) |
CRITICAL |
作业本身运行失败 |
工作流任务状态变化(flink:Workflow:TaskStateChange) |
INFO |
工作流任务状态变更通知 |
场景与效果
本示例以 Hologres 实时数仓搭建案例为背景,通过实时计算 Flink 版搭建 Hologres Streaming Warehouse 实时分层数仓,旨在通过DataWorks AI助手观测 Flink 流作业执行状态,以及通过多维度运行指标判定任务是否存在异常。
注:图片来源于阿里云官网文档。
实时数仓分层结构如下:
分层阶段 |
对应图上步骤 |
本示例对应 Flink 流作业节点 |
实现方式 |
ODS 层(业务数据库实时入仓) |
|
01toholo_database |
MySQL 业务表(如订单表、支付表、字典表)通过 Flink CDC(CDAS 语句)实时同步写入 Hologres,作为 ODS 层;同步时开启 Binlog,支持全量读取后自动切换为增量消费。 |
DWD 层(实时主题宽表构建) |
|
22consumodsbinlog_dwd_orders |
Flink 实时消费 ODS 层表,将多张源表进行维表 Join,打宽生成 DWD 层明细宽表并写回 Hologres。 |
DWS 层(实时指标聚合) |
|
33consumdwd_dws_shops、33consumdwd_dws_users |
Flink 消费 DWD 层宽表的 Binlog,实时聚合计算用户、商户维度指标,生成 DWS 层聚合表并写入 Hologres |
在实时计算场景下,数据时效性是关键衡量标准。数据新鲜度受业务数据输入、目标数据消费、任务代码逻辑以及 Flink 底层执行环境等多维度影响。本示例将对上述Flink 流作业节点运行健康度进行监控,预期效果如下:
Flink作业状态巡检结果示意 |
Flink作业指标巡检结果示意 |
Flink作业指标巡检HTML报告示意 |
|
|
|
步骤一:配置 Flink 作业失败事件巡检
本步骤通过 AI 助理配置定时巡检任务,监控 Flink 作业失败事件,出现 CRITICAL 事件时立即推送报警到钉钉群。
操作步骤:在钉钉群中向 AI 助理发送巡检任务创建指令,本示例指令如下:
帮我建个定时任务,命名为「Flink作业状态巡检」,每 1 分钟巡检一次 221812 空间的 4 个作业失败事件(JOB_FAILED),出现 CRITICAL 事件立即发送报告发到当前群并 @ 任务责任人:01toholo_database、22consumodsbinlog_dwd_orders、33consumdwd_dws_shops、33consumdwd_dws_users
注意:只有出现作业失败事件(JOB_FAILED),才需要推送到当前群。
请将指令中的工作空间 ID 和作业名称替换为您的真实值。
创建Flink作业状态巡检任务 |
Flink作业状态巡检结果示意 |
|
|
配置完成后,AI 助理会每分钟巡检一次 Flink 作业状态,当出现作业失败事件时,会自动推送报警到钉钉群并 @ 任务责任人。
步骤二:配置 Flink 作业指标巡检
本步骤通过 AI 助理配置定时巡检任务,对 Flink 作业的全部指标进行周期性巡检,并根据分级标准判定问题严重程度,输出结构化巡检报告。
指标组合判定
在实时计算场景下,数据时效性是关键衡量标准。数据新鲜度受业务数据输入、目标数据消费、任务代码逻辑以及 Flink 底层执行环境等多维度影响。同时,单看一个指标容易误判,实际排障需要看指标组合。以下是常见的指标组合及其判定结论可供参考:
指标组合 |
结论 |
分级 |
Failover↑ + Checkpoint=0 + Sink=0 |
崩溃循环,作业不可用 |
P0 · 致命(立即介入) |
GC 耗时↑ + 堆使用逼近 Max → Failover↑ |
OOM 型故障的完整前兆链 |
P1 · 严重(即将升级为 P0) |
输入正常 + 输出=0 + PendingRecords↑ |
反压/算子卡住 |
P1 · 严重(数据积压) |
业务延时↑ + 传输延时正常 |
瓶颈在计算 |
P2 · 警告(需扩容/调优) |
业务延时正常 + 传输延时↑ |
瓶颈在读取侧 |
P2 · 警告(需排查 Source/网络) |
输入=0 + 输出=0 + SourceIdleTime 高 |
上游没数据,预期空闲 |
P3 · 正常(无需处理) |
异常分级标准
为快速确认问题严重程度,本示例为不同指标组合设定异常级别。
级别 |
判定标准 |
处理建议 |
P0 — 作业已挂或数据已停 |
Failover ≥ 1,或 Checkpoint 持续为 0 达 5 分钟 |
立即处理 |
P1 — 数据延迟或链路异常 |
业务延时 ≥ 180 秒持续 3 周期;或 Sink 持续 0 达 5 周期且上游有输入;或积压持续上升 |
尽快处理 |
P2 — 资源水位偏高 |
TM CPU ≥ 85% 或堆使用率 ≥ 90% 持续 10 周期 |
关注 |
配置指标巡检
操作步骤:在钉钉群中向 AI 助理发送巡检任务创建指令,本示例指令如下:
帮我建个定时任务,命名为「Flink全量指标巡检」,每 5 分钟巡检 221812 空间绑定的 Flink 实例下 4 个作业的全部指标,报告发到当前群并@任务责任人:
01toholo_database、22consumodsbinlog_dwd_orders、33consumdwd_dws_shops、33consumdwd_dws_users
【分级标准】(对照判定,报告中标明级别):
- P0:Failover ≥ 1,或 Checkpoint 持续为 0 达 5 分钟
- P1:业务延时 ≥ 180 秒持续 3 周期;或 Sink 持续 0 达 5 周期且上游有输入(排除预期空闲);或积压持续上升
- P2:TM CPU ≥ 85% 或堆使用率 ≥ 90% 持续 10 周期
- P3:其他
【报告格式】:完整报告用HTML方式展示,摘要内容如下:
- 首行一句话结论:最高级别 + 异常作业数 + 与上轮相比的关键变化
- 异常项置顶:作业 / 指标 / 实测值 / 阈值 / 级别 / 持续时长;每个异常指标附一句解读(该指标正常时的意义 + 当前异常意味着什么)
- 结尾给组合研判:根因方向 + 下一步建议
【降噪】:与上轮相比级别和数值无变化的异常只汇总一行(标注「已持续 X 分钟」),不重复展开明细;级别变化或出现新异常时才输出完整解读
创建Flink作业指标巡检任务 |
Flink作业指标巡检结果示意 |
Flink作业指标巡检HTML报告示意 |
|
|
|
配置完成后,AI 助理会每 5 分钟巡检一次 Flink 作业的全部指标,自动判定问题严重程度,并输出结构化巡检报告到钉钉群。
场景扩展
以上案例覆盖了 Flink 作业本身的健康度巡检。在此基础上,您还可以进一步扩展:
- 上游业务库监控:针对 Flink 作业的上游业务库(如 RDS MySQL),配置数据库层面的监控报警。
- 下游数仓监控:针对 Flink 作业的下游数仓(如 Hologres),配置数据质量监控。
- 多通道推送:支持将巡检结果推送到钉钉、飞书、企业微信等多个 IM 通道。
加入官方交流群
您需要先单击申请链接加入"阿里云大数据AI平台"组织,再扫描下方二维码加入AI助理服务产品钉钉交流群,加入后,即可获得专属产品技术支持!
钉钉群号: 149605034971
钉钉群二维码: