数据库

使用 MongoDB 变更流热交换 Cron 调度

硬编码的 cron 意味着每当调度发生变化时都需要重新部署。将调度存储在数据库中,并让一个守护进程通过变更流来响应编辑。以下是该设计及其故障处理。

本文由 AI 模型从英文原文翻译而来,措辞可能与原文有出入。 阅读英文原文

如果你的后台作业基于硬编码在代码或环境变量中的 cron 表达式运行,那么更改作业的运行时间就意味着需要重新部署。将一个报告的生成时间从上午 9 点改到上午 8 点,你就需要交付一个构建版本、运行一个流水线并重启一个进程,而这一切只为了修改一个数字。这个过程很慢,与变更本身相比风险过高,并且不鼓励进行那些能保持系统健康的微小调整。解决方案是将调度计划存储在数据库中,并让调度器实时响应编辑操作。MongoDB 变更流使这种响应能够即时发生。本文将介绍其设计、使其可靠的重连处理机制,以及在何种情况下它会比你实际需要的机制更复杂。

代码中硬编码调度的问题

硬编码的调度将两件变化速率和变化原因完全不同的事物耦合在了一起。作业的逻辑随工作内容的变化而变化,这是一个属于构建(build)范畴的工程决策。作业的执行时间随运营需求的变化而变化——例如,报告应在会议前生成,重型任务应移出高峰时段——而这是一个与代码无关的运营决策。将它们绑定在一起,意味着每一次对执行时间的微调,都拖累着一次完整的代码部署。

其代价不仅仅是部署时间,更是其带来的寒蝉效应。当更改调度需要一次生产环境的部署时,人们就会停止这样做。调度会逐渐偏离实际需求,因为没人愿意为了将一个作业的执行时间移动一小时而去发布一个构建版本。调度应该是你可以编辑的数据,而不是你必须发布的代码

将调度视为数据

第一步是将调度放入一个集合中,每个作业一个文档,并将 cron 表达式作为一个字段。

// One document per scheduled job. The cron expression is data, not code.
type Schedule struct {
	ID        string    `bson:"_id"`        // e.g. "daily_report"
	CronExpr  string    `bson:"cron_expr"`  // e.g. "0 9 * * *"
	UpdatedAt time.Time `bson:"updated_at"`
}

现在,更改调度只需一次数据库更新,无需部署:

db.schedules.updateOne(
  { _id: "daily_report" },
  { $set: { cron_expr: "0 8 * * *", updated_at: new Date() } }
)

但仅有数据是不够的。需要有东西注意到变更并重新配置正在运行的调度器。你有两种方法可以实现这一点,而它们之间的区别正是本文的重点。

轮询与响应

简单的方法是轮询:每分钟,调度器会重新读取集合,并在有任何更改时重建其计时器。这种方法是可行的,并且对于许多系统来说完全足够。它的缺点是延迟和浪费。一项更改最多需要一个完整的轮询间隔才能生效,而且绝大多数轮询都发现没有任何变化,纯属开销。缩短间隔以减少延迟,就会增加浪费。它就像一个没有好设置的刻度盘,只有不那么坏的设置。

变更流移除了这个刻度盘。调度器不再按计时器询问“有变化吗?”,而是订阅一个实时源,并在发生变化时立即得到通知。一次计划编辑会近乎实时地传播,当没有变化时,调度器不执行任何操作,也不产生任何成本。你同时获得了更低的延迟和更少的工作量,这非同寻常,值得为该机制付出额外的精力。

什么是变更流

变更流是一个实时、有序的数据源,其中包含对集合执行的写入操作。在底层,它读取数据库用于保持副本集成员同步的同一份复制日志,因此它能按顺序看到已提交的 inserts、updates 和 deletes,无需轮询。您的守护进程只需打开一次流,然后就会阻塞,每当计划文档发生变更时,它都会收到一个事件。

// Watch the schedules collection. Block until an edit arrives, then apply it.
stream, err := coll.Watch(ctx, mongo.Pipeline{})
if err != nil {
	return fmt.Errorf("open change stream: %w", err)
}
defer stream.Close(ctx)

for stream.Next(ctx) {
	var event struct {
		FullDocument Schedule `bson:"fullDocument"`
	}
	if err := stream.Decode(&event); err != nil {
		logger.Warn("decode change event", logger.Err(err))
		continue
	}
	scheduler.Reload(event.FullDocument) // rebuild this job's timer
}

当操作员运行那一行更新时,守护进程会立即收到事件,并仅重建受影响的计时器。无需重新部署,无需重启,也无需等待轮询间隔。现在,计划是真正的实时数据。

变更流需要副本集,因为它们读取的是复制日志,而独立服务器不会维护该日志。任何生产环境的 MongoDB 部署都已经是副本集,因此这很少成为一项新要求,但在本地单节点设置中使用此功能之前,您有必要了解这一点,因为在该环境下此功能不可用。

使其可靠的部分:恢复令牌

一个简单的变更流存在间隙。如果连接因网络抖动、故障转移或短暂重启而中断,而你只是简单地重新打开流,那么你会从当前时间点恢复,并静默地错过在断开连接期间发生的所有变更。在三十秒的重新连接期间编辑的计划将会丢失,守护进程会继续按旧的时间运行,而数据库中的数据却已不同。这种不一致正是整个设计旨在防止的错误。

变更流通过恢复令牌解决了这个问题。每个事件都携带一个令牌,标记其在信息流中的位置。持久化你已处理的最新令牌,在断开连接后,从该令牌而不是从当前时间点开始重新打开流。数据库会按顺序重放你错过的每一个变更,让你能准确地从中断处继续。

// Resume from the last processed position so a reconnect misses nothing.
opts := options.ChangeStream()
if resumeToken != nil {
	opts.SetResumeAfter(resumeToken)
}
stream, err := coll.Watch(ctx, mongo.Pipeline{}, opts)
// ... after handling each event, persist stream.ResumeToken() durably

关键在于,在处理完每个事件后都要存储令牌,并且要足够持久以应对守护进程的重启。重新连接时,加载该令牌并从其位置恢复。有了这一机制,变更流不仅速度快,而且在任何长期运行的消费者都不可避免的断连情况下也能做到无间隙。正是令牌将一个便捷的功能转变为一个可靠的功能。

一个注意事项:恢复令牌仅在变更仍然存在于复制日志中时才有效,而复制日志的大小是有限的。如果守护进程停机时间过长,导致最旧的未处理变更已从日志中清除,那么令牌将变得无效,流也无法从该令牌恢复。处理这种情况的方法是,检测无效恢复错误,然后回退到对集合进行完全重载,这将调度器重新同步到当前真实状态。这与你无论如何都希望在首次启动时执行的从头开始协调的路径是相同的。

平滑应用变更

接收事件只是工作的一半。应用变更时必须小心,因为你正在修改一个正在运行的调度器。应专门为已变更的作业重建计时器,而不是拆除并重新创建每个计时器,这样对一个调度计划的编辑就不会干扰其他计划,或在切换过程中冒着重复触发的风险。取消该作业 ID 的旧计时器,根据更新后的 cron 表达式安装新的计时器,并保持所有其他作业不受影响。保持重新配置的幂等性,这样即使同一个事件被应用两次(这在重连边界上可能由 resume-from-token 引起),调度器也能达到与应用一次相同的状态。正是这种幂等性,使得恢复的数据流的“至少一次”特性变得安全。

何时应保持简单

当调度变更足够频繁,以至于重新部署成为真正的拖累时,或者当变更生效的延迟确实很重要时,变更流就是合适的工具。如果你的调度几乎从不更改,那么在运维上带来的好处就很小,并且简单的“重启以重新加载”或慢轮询的方式,需要构建的东西更少,也更不容易出问题。不要为一个一年只编辑两次调度的系统添加实时事件消费者、恢复令牌持久化以及无效令牌回退机制。

更深层次的原则不仅仅适用于 cron。任何出于运维原因而非代码原因而更改的配置,都应该存放在一个供运行中的系统监视的数据存储中,这样一来,更改它就只是一个编辑操作,而不是一次发布。调度就是一个清晰的例子,因为其带来的好处非常具体:过去需要通过部署才能完成的时间变更,现在变成了一行更新。但同样的模式也适用于功能开关、路由规则和速率限制。将配置作为数据存储,使用变更流来监视它,并使用恢复令牌来处理重连。这三者的组合正是使实时配置值得信赖而不仅仅是方便的原因。