实时工作流的异常处理

更新时间:
复制 MD 格式

实时工作流的异常处理分为算子级、事件级和工作流级三个层级,本文为您详细介绍异常处理过程。

算子级异常处理

image

策略

说明

适用场景

自动重试

算子执行失败后自动重试,最大重试次数可配置(默认为3次),重试间隔支持固定间隔或指数退避。

临时性错误(网络超时、容器短暂不可用)。

超时控制

单个算子执行超过配置的超时时间(默认为300秒)后,强制终止并标记为失败。

算子逻辑阻塞、大文件处理超时等。

事件级异常处理

数据集写入失败处理机制:
若算子处理成功但写入目标数据集失败(例如连接异常、Schema不匹配、主键冲突等),将遵循事件级异常处理策略,即自动将异常数据归档至脏数据文件并记录异常日志,此时工作流继续运行,同时提交当前位点并继续消费后续事件。写入失败的具体错误信息可前往研发 > 任务运维 > 实例运维 > 实时实例中,在对应实时实例的运行日志中查看脏数据文件。

事件级异常产出物:

  • 脏数据文件:写入CSV(产生时间、报错内容、报错原因),详情请参见。

  • 告警触发:若失败频率超过配置阈值,触发失败频率超过配置告警。

工作流级异常处理

异常场景

系统行为

恢复方式

Listener崩溃

Session监控服务检测到Listener心跳丢失,自动重启Listener。

从最近Checkpoint恢复消费位点,uncommit set中的事件自动重放。

Credit余额不足

工作流任务不会终止,但会终止对Kafka的消费。

充值后,自动恢复(从最近Checkpoint恢复消费位点)。

算子容器OOM/崩溃

Kubernetes自动重启容器(RestartPolicy=Always)。

重启后Router重新纳入该容器;处理中的事件由编排层超时后重试到其他容器。

对象存储不可用

Context溢出写入失败,算子执行异常。

将持续触发算子级重试,若最终无法处理,则视为脏数据。

Kafka连接断开

Listener自动重连(Kafka consumer内置重连机制)。

重连成功后从当前位点继续消费,若超时未恢复,则Listener标记为异常并重启。

Checkpoint与恢复

说明

Checkpoint内容

仅保留start位点(序列化),start-end窗口长度为1000。

Checkpoint频率

内置策略不可配置,每60秒执行一次。

Checkpoint存储

持久化到数据库。

恢复模式

断点续跑:从最近一次成功的Checkpoint恢复start/end/uncommit set,uncommit set中的事件会被重新消费处理。

启动模式

  • 全新启动:从指定消费起点(最早/最新/指定时间位点)开始消费。

  • 恢复启动:从最近Checkpoint断点续跑。

幂等保障

由于Checkpoint恢复可能导致事件重复消费,算子和数据集写入需具备幂等性(如主键upsert);非幂等场景需由业务侧保障去重。