破坏性读取将解析失败转变为永久性挂起
某个 worker 使用 GETDEL 读取任务结果,因此任何它无法解析的负载都会被永久删除。将此情况视为“尚未就绪”意味着永久等待。
一个作业在 18:28 失败了。用户在 20:30 之前一直看到“进行中”的状态。失败通知已准时送达,工作单元也已读取。只是它没有理解通知的内容。
这个 bug 是一行控制流代码的问题。一个等待带外结果的工作单元将所有获取错误都同等对待,因此“结果还没准备好”和“结果显示作业已失败”都意味着继续等待。单是这一点,就只会造成一个缓慢的 bug。但读取操作本身使问题变成了永久性的:该工作单元使用了破坏性读取,因此在它误读失败通知的那一刻,该通知就不复存在了。
如果你有一个异步工作单元通过轮询一个键来获取结果,那么本文将探讨把一个小小的分类错误变成数小时挂起的两个因素,以及如何构建消费者来避免这种情况发生。
消费者必须区分的两种状态
这种设置很常见。后端向计算服务提交工作,并立即获得一个 ID。计算服务执行其任务,将结果写入一个共享键,然后发布一个“轻推”通知。后端订阅这个“轻推”通知,并且还通过定时器轮询该键,因为发布/订阅的投递不保证成功。
消费者循环大致如下:
for {
select {
case <-ctx.Done():
return rescueOrError(ctx.Err())
case msg := <-notifyCh:
result, err := fetchResult(ctx, key)
if err != nil {
slog.Warn("fetch after notify failed", "error", err)
continue // keep waiting
}
return result, nil
case <-ticker.C:
result, err := fetchResult(ctx, key)
if err == nil {
return result, nil
}
// no result yet, keep waiting
}
}
再读一下那个 continue。当计算服务明确告知我们作业失败时,就会执行到它。消费者会记录一条警告并返回休眠状态。
一个等待中的消费者必须区分至少三种结果,而这个循环将它们合并成了两种:
- 尚未就绪。 键不存在。继续等待。这是唯一应该继续等待的情况。
- 暂时性的基础设施错误。 键存储暂时无法访问或正在加载。继续等待,因为结果可能还未被消费。
- 终结状态。 结果已到达,它要么是明确的失败,要么是我们无法解释的东西。立即停止。
修复方法不是“处理错误情况”。而是将第三类情况变成循环可以识别的一种类型,并将每个分支都通过同一个判定条件进行路由:
type StepError struct{ Step, Server, Message string } // service said: failed
type ProtocolError struct{ Step, Server, Reason, Raw string } // we cannot parse it
func isTerminal(err error) bool {
var se *StepError
if errors.As(err, &se) { return true }
var pe *ProtocolError
return errors.As(err, &pe)
}
有两个细节比表面上看起来更重要。将它们作为指针返回,因为 errors.As 针对 *StepError 目标时,会静默地匹配值类型失败。并且使用 %w 进行包装,绝不使用 %v,否则链将在第一个包装器处断开,并且每个下游检查都会静默地返回 false。
为什么破坏性读取会增加风险
消费者通过原子性的读取并删除操作(Redis GETDEL)来读取结果。这个选择本身是合理的:它能防止两个工作单元消费同一结果,并且无需依赖 TTL 即可保持键空间整洁。
但这也将投递保证转换为了至多一次,而这就改变了解析失败所代表的含义。
对于一个普通的 GET 请求,格式错误的负载会很烦人。你可以记录它,将值保留在原位,然后稍后再做决定。而对于 GETDEL,值在查看它的那一刻就被消费掉了。因此,规则变得绝对:
一旦破坏性读取返回了除“键缺失”之外的任何内容,其后的每一步失败都是终结性的。
无效的 JSON 是终结性的。缺失 status 字段是终结性的。status 字段是数字而非字符串也是终结性的。这些问题都无法通过等待来解决,因为已经没有什么可等待的了。最初的循环将所有这些情况都视为“继续等待”,因此,每种情况都保证会导致数小时的挂起,直到某个外部安全计时器触发。
这里存在一个真正的权衡。如果你保留破坏性读取,就得接受消费者在读取和提交之间崩溃会导致结果丢失。如果你切换到 GET 加上在自身状态提交后执行删除的模式,你就能实现幂等性,但必须处理重复消费的问题。两种方式都可以。不可以的是,使用破坏性读取的同时,编写的消费者代码却假定它可以再次查看。
肯定式白名单,而非否定式检查
在添加终端类型时,我发现同一个函数中隐藏着第二个 bug,它比我正在修复的那个更旧:
var status string
if s, ok := raw["status"]; ok {
json.Unmarshal(s, &status) // error ignored
}
if status == "error" {
return nil, errors.New("service reported failure")
}
return &Result{Status: "done", Content: raw}, nil // everything else
这会检查那一个坏值,并将其余所有情况都视为成功。我们来逐步分析这意味着什么:
| 有效负载 | 旧行为 | 正确行为 |
|---|---|---|
{"status":"done", ...} |
完成 | 完成 |
{"status":"error", ...} |
失败,然后被调用方忽略 | 终结性失败 |
{} |
完成 | 终结性协议错误 |
{"status":123} |
完成 (反序列化错误被丢弃) | 终结性协议错误 |
{"status":"failed"} |
完成 | 终结性协议错误 |
| 完全不是 JSON | 错误,然后被调用方忽略 | 终结性协议错误 |
其中有三行是静默的数据损坏。一个空对象被报告为作业成功,其附带的任何部分内容都因此被写入记录。没有人注意到这一点,因为生产者恰好总是发送格式良好的有效负载。契约的履行是靠运气,而不是靠强制执行。
反过来。明确指定你接受的值,并将其他所有情况都视为你能看到的错误:
switch status {
case "done":
return &Result{Status: "done", Content: raw}, nil
case "error":
return nil, &StepError{Step: step, Message: sanitize(msg)}
default:
return nil, &ProtocolError{Step: step, Reason: "unknown_status:" + status}
}
同样的准则也适用于你将要忽略的解组错误。如果 status 存在但不是字符串,那么这就是一种生产者契约违规,而且你希望它能大声报错,而不是被转换成虚假的成功。
重试预算需要其自身的键作用域
一旦故障立即显现,你可能就想要重试。计算端的瞬时内存不足错误值得再尝试一次,而错误的输入则不值得,但你通常无法从有效负载中区分它们。
我的第一次尝试复用了一个现有的“我们已尝试过的服务器”数组作为预算计数器。两位评审者相继否决了它,而第二个原因很有意思。
第一个问题是差一错误 (off-by-one)。追加当前服务器然后检查 len(servers) >= 1 会在第一次尝试时就导致作业失败,因此重试永远不会发生。一说出来就很明显了。
第二个问题是当你修复第一个问题后会发生什么。将上限提高到 2,每当同一个工作进程重新接手该作业时,数组就会停止增长,因为追加操作有防止重复的保护。没有增长就意味着没有终止条件。作业会永远重新排队。更糟糕的是,通常的应急出口(escape hatches)缺失了:这条路径没有增加通用的重试计数器,并且作业在每个周期都干净地完成,因此过时作业监视器从未发现它被卡住。
同样的结构之所以对邻近的一类错误是安全的,原因在于熔断器。提交失败会被计入其中,所以几次失败后熔断器就会打开,工作进程进入空闲状态,这意外地打破了循环。我所做的更改明确地移除了这个意外。
所以,预算需要是一个单调递增的计数器,并且其键作用域需要仔细考虑:
// key: retry:{job_id}:{step}
const incrWithTTL = `
local v = redis.call("INCR", KEYS[1])
if v == 1 then redis.call("EXPIRE", KEYS[1], ARGV[1]) end
return v`
在最终确定此方案前,我犯了三个错误:
- 按作业划分作用域,而非按实体。 基于父实体加上步骤名称来生成键(key)看起来很自然,但如果一个用户可以针对同一个实体和步骤触发多个独立的作业,这些作业会共享同一个预算并相互耗尽资源。每个入队操作使用一个独立的 ID 可以将它们分开,而内部的重新入队操作会保留这个 ID。
- 分离子步骤。 在同一个作业记录中运行的预处理步骤,绝不能消耗主步骤的预算。使用实际的计算步骤名称,而不是作业声明的步骤名称。
- 让 INCR 和 EXPIRE 原子化。
INCR创建键时不会设置 TTL。如果进程在执行EXPIRE之前死掉,你就会得到一个永不过期的计数器,此后该 ID 的每一次失败都会立即致命,且无法自然恢复。使用一个 Lua 脚本可以消除这个时间窗口。
应根据两次尝试之间可能的最长间隔来设置 TTL,而不是根据平均间隔。我最终设置的值是单次等待硬性上限的两倍。我最初的两次猜测(一小时,然后是六小时)都比单次尝试可能运行的时间要短,这会导致计数器在重试过程中过期并重置预算。这就又变成了无限重新入队问题,只是换了个形式。
不要将领域故障混入你的熔断器
最后一部分是关于这些故障是否应该影响生产者的健康状况。
熔断器的存在是为了回答一个问题:这个端点是否生病了?一个明确的失败响应是它健康的证据。它接受了请求,完成了工作,生成了结构化的答案,并通过结果通道交付了它。是任务失败了,而不是服务器失败了。
因此,明确的失败不应算作熔断器故障。否则,少数几个错误的输入就会打开熔断器,将一个完全正常的节点移出轮换,这会将不相关的工作推入重新排队的队列中。
协议错误则相反。一个你无法解析的负载表明可能存在版本偏差、序列化错误或未完成的部署。这是一个节点问题,应该被计算在内。
这里有一个微妙之处,它让我多花了一次提交来修复。如果你只是跳过失败调用,你可能也会跳过成功调用,并且如果你的熔断器在半开状态下跟踪一个进行中的探测槽,那么该槽将永远不会被释放。然后,工作单元会停留在一个永远不会打开的门后,直到进程重启。我需要一个明确的释放操作:
func (cb *Breaker) ReleaseProbeAsSuccess() {
cb.mu.Lock()
defer cb.mu.Unlock()
if cb.state == "half-open" {
cb.state = "closed" // the server answered, it is reachable
cb.failures = 0
}
cb.inFlight = 0
}
该方法有两项防护措施。不要从断路器关闭期间发起的请求中调用它,因为断路器通常在 goroutine 之间共享,你可能会在没有真正的探测成功的情况下关闭别人的探测。并且,绝不要在上下文取消时调用它,因为一次关闭操作并不能证明远程主机的任何情况。
仍未涵盖的内容
坦诚面对边界情况,因为它们是下一个错误的来源:
- 一个作业的并发消费者。 如果两个工作单元以某种方式处理了同一个作业,其中一个可能会在计数为 1 时重新入队,而另一个在计数为 2 时放弃,导致一个已经失败的作业返回给生产者。要防范这种情况,需要在重新入队时对作业状态进行比较并设置(compare-and-set),而不仅仅是使用一个计数器。
- 缺乏多样性的重试。 计数器限制了总尝试次数,但并未将重试引导至不同的节点。如果故障是节点本地的,盲目重试会浪费重试预算。
- 生产者合约本身。 所有这些都是消费者在防御模糊性。真正的修复方法是让有效负载(payload)明确说明故障是否可重试,这样消费者就不必猜测。
部署后,同一个失败的作业在两个节点上尝试了两次,并在 2 分 19 秒内处理完毕,生产者的原始错误文本被保留在记录中。而之前的路径耗时超过六个小时,并且用一个通用的超时标签覆盖了失败原因。
常见问题解答
破坏性读取是错误的吗?
不是。这是一种保证单次消费的合理方式。它只是转移了责任:一旦消费,解析就是你最后的机会,因此每次读取后的失败都必须是终结性的。错误在于将破坏性读取与假定可以二次查看的消费者代码配对。
为什么不在解析错误时直接重试抓取?
因为对于破坏性读取,已经没有什么可抓取的了。重试操作看到的是一个空键,结论是“尚未就绪”,然后等待。这正是解析错误如何演变成长达数小时的挂起的原因。
终结性失败应该尝试多少次?
重试一次对我来说是合适的,因为有相当一部分失败是计算端的瞬时资源耗尽,而最坏情况下的成本只是额外运行一次该步骤。如果你的生产者报告某个错误是否可重试,就使用该信息,而不是一个固定的数字。固定数字是你所没有的信息的替代品。
协议错误和显式失败都应该重试吗?
我两者都重试,因为廉价的统一策略比我无法验证的分类法更容易推理。它们的区别在于如何影响熔断器,而不在于获得多少次尝试。
如何安全地存储生产者的错误文本?
在截断之前进行规范化,而不是之后。修复无效的 UTF-8,剥离控制字符和双向覆盖符,隐去内部路径和 URL,然后裁剪到指定长度。先裁剪会留下半个路径或 URL,使其不再匹配你的隐去模式。并且,不要让原始未解析的负载进入任何会触达用户的字段,因为它是整个系统中最不可信的字符串。