实时工作流的异常处理分为算子级、事件级和工作流级三个层级,本文为您详细介绍异常处理过程。
算子级异常处理
策略 | 说明 | 适用场景 |
自动重试 | 算子执行失败后自动重试,最大重试次数可配置(默认为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内容 | 仅保留 |
Checkpoint频率 | 内置策略不可配置,每60秒执行一次。 |
Checkpoint存储 | 持久化到数据库。 |
恢复模式 | 断点续跑:从最近一次成功的Checkpoint恢复 |
启动模式 |
|
幂等保障 | 由于Checkpoint恢复可能导致事件重复消费,算子和数据集写入需具备幂等性(如主键upsert);非幂等场景需由业务侧保障去重。 |