ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

Bytebase 任务运行高可用改造:用 TaskRunPayload JSONB 列持久化 SchedulerInfo

Bytebase 任务运行高可用改造:用 TaskRunPayload JSONB 列持久化 SchedulerInfo Bytebase 任务运行高可用改造用 TaskRunPayload JSONB 列持久化 SchedulerInfo【免费下载链接】bytebaseDatabase governance built for humans and agents — controlling changes and access across every major database.项目地址: https://gitcode.com/GitHub_Trending/by/bytebase本设计文档记录了 Bytebase 后端一次关键的高可用HA架构改造将原本保存在进程内存Bus.sync.Map中的任务调度信息SchedulerInfo迁移到task_run表的payloadJSONB 列中持久化存储从而让调度状态在多副本之间共享、在副本崩溃后不丢失并让 UI 在任何副本存活的场景下都能展示任务等待原因。读完本文你将掌握这条Proto 定义 → 数据库迁移 → Store 读写 → 调度器写入 → API 输出的完整改造链路以及 Bytebase 在 HA 模式下任务调度状态管理的底层实现方式。问题背景内存态调度信息如何破坏高可用Bytebase 的任务调度器Task Scheduler负责将PENDING状态的任务运行逐步推进到AVAILABLE、RUNNING并最终完成。在这一过程中调度器需要记录某个任务运行当前为什么在等待例如因为达到并行任务数上限而被阻塞。改造之前这部分信息TaskRunSchedulerInfo被保存在Bus.sync.Map中——即仅存在于单个副本的进程内存里。从设计文档 2026-01-13-task-run-payload-ha-design.md 可以看到这种方案在高可用部署下存在三个明确缺陷副本崩溃即丢失保存调度信息的副本进程一旦崩溃内存中的所有调度状态随之消失副本之间不共享调度信息只写入了执行调度的那个副本其他副本无法访问UI 无法展示等待原因当负责追踪的副本死亡后前端界面无法向用户解释任务为什么还在等待。在 Bytebase 的高可用架构中多个副本通过 advisory lock集群级互斥锁竞争调度权任一副本都可能随时接管调度职责。内存态状态显然无法在副本间平滑交接这正是本次改造要解决的核心问题。解决方案为 task_run 表引入 payload JSONB 列设计文档给出的方案非常直接在task_run表中新增一个payloadJSONB 列用来存放一个TaskRunPayloadproto 消息其中包含SchedulerInfo。该方案的关键优势在于持久化调度信息写入数据库不再依赖任何单个副本的存活多副本共享任何副本都能从数据库读取到最新的调度状态schema 可演进未来需要在 payload 中增加字段时只需扩展 proto 消息无需再执行数据库 schema 变更。这条思路与 Bytebase 其他模块如review_config、workspace的 payload 迁移保持一致属于项目中proto JSONB持久化模式的又一次复用。Proto 设计可扩展的 TaskRunPayload 消息改造的第一步是在存储层 proto 中定义新消息。设计文档规划了TaskRunPayload的结构而当前仓库 proto/store/store/task_run.proto 中的实现已经落地并有所演进// TaskRunPayload contains extensible runtime data for a task run. // Stored in the payload JSONB column. New fields can be added here // without database schema changes. message TaskRunPayload { // Scheduler information about why a task is waiting. SchedulerInfo scheduler_info 1; // If true, prior backup is skipped for this task run. bool skip_prior_backup 2; }注意其中的skip_prior_backup字段它在设计文档之后被追加恰好实证了未来字段可以加入 proto 而无需 schema 变更这一设计承诺——payload 列依然保持原有结构仅靠 proto 扩展就承载了新的业务语义。SchedulerInfo消息本身包含两个字段message SchedulerInfo { // Timestamp when the scheduler reported this information. google.protobuf.Timestamp report_time 1; // WaitingCause indicates why a task run is waiting to execute. message WaitingCause { reserved 1 to 2; oneof cause { // Task is waiting due to parallel execution limit. bool parallel_tasks_limit 3; } } // Reason why the task run is currently waiting. WaitingCause waiting_cause 2; }waiting_cause采用oneof结构目前唯一的原因是parallel_tasks_limit并行任务数限制report_time记录调度器报告该信息的时间戳。reserved 1 to 2说明该 oneof 曾经历过字段演进为后续新增等待原因如数据库互斥等待预留了扩展空间。值得说明的是同一份SchedulerInfo结构也以OUTPUT_ONLY字段的形式暴露在 v1 API 的TaskRun消息中见 proto/v1/v1/rollout_service.proto 第 545-562 行字段编号为 12包含report_time与waiting_cause用于向客户端输出调度信息。Schema 迁移一行动态添加 payload 列数据库迁移文件 0032##task_run_payload.sql 与设计文档完全一致只有一条语句ALTER TABLE task_run ADD COLUMN payload jsonb NOT NULL DEFAULT {};三个关键点jsonb类型PostgreSQL 二进制 JSON 格式支持高效的查询与索引NOT NULL保证每一行都有 payload 值简化代码中空值判断DEFAULT {}存量行自动获得空对象作为初始 payload保证迁移对历史数据零侵入、可平滑上线。从迁移文件命名3.14/0032##task_run_payload.sql可以看到该改造随 3.14 版本一起发布。按照设计文档的清单同步还需要更新backend/migrator/migration/LATEST.sql追加该列以及backend/migrator/migrator_test.go更新版本号这是 Bytebase 迁移体系的标准流程。Store 层TaskRunMessage 的读写闭环消息结构扩展在 backend/store/task_run.go 中TaskRunMessage结构体新增了PayloadProto字段第 24 行type TaskRunMessage struct { TaskUID int64 Environment string // Refer to the tasks environment. PlanUID int64 Status storepb.TaskRun_Status ResultProto *storepb.TaskRunResult PayloadProto *storepb.TaskRunPayload ... }注意这里同时保留了ResultProto同样以 JSONB 形式存储的TaskRunResult说明 payload 列的设计与既有的 result 列一脉相承都采用proto → JSONB的持久化模式。读取路径SELECT protojson 反序列化ListTaskRuns的 SQL 查询将task_run.payload加入 SELECT 列表第 114 行并在行扫描后用protojson.Unmarshal将其反序列化为TaskRunPayloadproto第 190-194 行var payloadProto storepb.TaskRunPayload if err : common.ProtojsonUnmarshaler.Unmarshal([]byte(payloadJSON), payloadProto); err ! nil { return nil, errors.Wrapf(err, failed to unmarshal task run payload: %s, payloadJSON) } taskRun.PayloadProto payloadProto写入路径UpdateTaskRunPayload 方法Store 层提供了专用写入方法UpdateTaskRunPayload第 670-693 行将 payload proto 序列化后按(projectID, taskRunID)更新func (s *Store) UpdateTaskRunPayload(ctx context.Context, projectID string, taskRunID int64, payload *storepb.TaskRunPayload) error { payloadBytes, err : protojson.Marshal(payload) ... payloadStr : string(payloadBytes) if payloadStr { payloadStr {} } ... SET payload ?, updated_at now() ... }空 proto 序列化结果会被归一化为{}与表定义的默认值保持一致。此外在批量创建 task run 时第 348-356 行创建方传入的PayloadProto也会一并序列化写入确保创建即携带初始 payload。调度器在哪里写入、何时清除调度主循环与集群互斥调度器位于 backend/runner/taskrun/pending_scheduler.go。runPendingTaskRunsScheduler以 5 秒为周期taskSchedulerInterval定义于 scheduler.go运行同时响应TaskRunPendingTickleChan事件触发。每一轮调度先通过TryAdvisoryXactLock(ctx, tx, store.AdvisoryLockKeyPendingScheduler)获取集群级互斥锁保证任意时刻只有一个副本在推进 PENDING 任务这是 HA 模式下谁调度谁写状态的关键前提——因为只有一个调度者写库的调度信息天然一致。写入等待原因storeParallelLimitCause设计文档要求storeParallelLimitCause()将TaskRunPayload{SchedulerInfo: ...}写入数据库。当前实现第 194-208 行完全遵循该设计func (s *Scheduler) storeParallelLimitCause(ctx context.Context, projectID string, taskRunID int64) { payload : storepb.TaskRunPayload{ SchedulerInfo: storepb.SchedulerInfo{ ReportTime: timestamppb.Now(), WaitingCause: storepb.SchedulerInfo_WaitingCause{ Cause: storepb.SchedulerInfo_WaitingCause_ParallelTasksLimit{ ParallelTasksLimit: true, }, }, }, } if err : s.store.UpdateTaskRunPayload(ctx, projectID, taskRunID, payload); err ! nil { slog.Error(failed to store parallel limit cause, log.BBError(err)) } }它发生在schedulePendingTaskRun的并行限制检查Check 5失败时当sc.checkParallelLimit(task, maxParallel)判定当前活动任务数已达到上限maxParallel来自project.Setting.GetParallelTasksPerRollout()即项目设置的每个 rollout 并行任务数调度器不直接推进该任务而是先记录等待原因再返回等待下一轮调度周期再次尝试。清除等待原因promoteTaskRun当所有前置检查通过、任务被提升为AVAILABLE时promoteTaskRun第 230-256 行会先写入空 payload 清除旧的调度信息// Clear scheduler info by writing empty payload if err : s.store.UpdateTaskRunPayload(ctx, taskRun.ProjectID, taskRun.ID, storepb.TaskRunPayload{}); err ! nil { slog.Error(failed to clear scheduler info, log.BBError(err)) }随后才执行状态更新PENDING → AVAILABLE带AllowedStatuses做乐观并发控制并通过TaskRunRunningTickleChan唤醒运行期调度器。写入空 payload 而非删除列值与NOT NULL DEFAULT {}的 schema 定义相契合——没有等待原因本身就是一种明确的、持久化的状态。完整的调度检查链将 payload 写入点放回调度流程中可以看到全貌。schedulePendingTaskRun依次执行详见 pending_scheduler.go 第 135-192 行Check 1RunAt时间未到则不调度项目检查项目被删除则跳过实例检查目标实例被归档则跳过Check 4数据库互斥检查顺序任务在同一数据库上互斥Check 5并行任务限制——失败时写入SchedulerInfo{WaitingCause: parallel_tasks_limit}全部通过后promoteTaskRun先清空 payload 再置为AVAILABLE。调度信息的写入等待时与清除提升时就这样嵌入了状态机流转的两个端点任何一个副本执行调度状态都会落库。API 层从 PayloadProto 读取 SchedulerInfo设计文档要求 backend/api/v1/rollout_service_converter.go 移除bus *bus.Bus参数、改为从taskRun.PayloadProto.SchedulerInfo读取。当前实现正是如此if taskRun.PayloadProto ! nil taskRun.PayloadProto.SchedulerInfo ! nil { t.SchedulerInfo convertToSchedulerInfo(taskRun.PayloadProto.SchedulerInfo) }配套的转换函数将 store 层 proto 映射为 v1 API protoconvertToSchedulerInfo拷贝report_time与waiting_cause第 59-68 行convertToSchedulerInfoWaitingCause通过 type switch 识别ParallelTasksLimit分支转换为 v1 的parallel_tasks_limit字段第 70-84 行未知分支返回nil保证前向兼容。由此任何调用GetTaskRun/ListTaskRunsAPI 的客户端包括前端 UI都能直接拿到scheduler_info字段且该字段在 v1 proto 中被标记为OUTPUT_ONLY仅服务端可写、客户端只读。UI 展示任务因并行限制而等待时数据来源从某个副本的内存变成了数据库中的持久化记录彻底消除了副本存活依赖。清理与落地清单设计文档的最后一步是清理内存态存储。从当前仓库搜索来看TaskRunSchedulerInfo已从 backend/component/bus/bus.go 中移除仓库中已无该符号sync.Map内存缓存方案被完全废弃任务调度信息的内存态生命周期正式结束。设计文档列出的完整改动清单9 项如下全部已落地文件改动内容当前状态proto/store/store/task_run.proto新增TaskRunPayload含SchedulerInfo✅ 已实现并追加skip_prior_backup字段backend/migrator/migration/3.14/0032##task_run_payload.sql新增 payload 列新文件✅ 存在backend/migrator/migration/LATEST.sql追加该列✅ 随版本合并backend/migrator/migrator_test.go更新版本号✅ 随版本合并backend/store/task_run.go新增PayloadProto字段与读写方法✅ 已实现backend/runner/taskrun/pending_scheduler.gostoreParallelLimitCause/promoteTaskRun使用 store✅ 已实现backend/api/v1/rollout_service_converter.go移除 Bus、改读PayloadProto.SchedulerInfo✅ 已实现backend/api/v1/rollout_service.go更新 converter 调用✅ 已实现backend/component/bus/bus.go移除TaskRunSchedulerInfo sync.Map✅ 已移除从仓库现状可以推断该设计文档对应的改造已经完整落地并随版本发布store 层的读写闭环、调度器的写入/清除点、API 层的输出转换全部就位且TaskRunPayload已经承载了设计文档之外的第二个业务字段。这也说明这套proto JSONB的持久化模式在 Bytebase 内部具有良好的可扩展性——后续任何需要随任务运行持久化的运行时元数据都可以继续在TaskRunPayload中追加字段而无需再次触碰数据库 schema。总结本次改造的核心价值可以用三句话概括调度状态从进程内存走向数据库持久化从单副本可见走向多副本共享从崩溃即丢失走向随时可恢复。对于 Bytebase 的用户而言最直接的收益是在高可用部署下UI 中的任务等待原因不再受限于某个副本的存亡对于开发者而言TaskRunPayload提供了未来扩展运行时元数据的低成本通道。这条问题分析 → Proto 设计 → 迁移 → Store → Scheduler → API → 清理的改造路径也完整展示了 Bytebase 如何在一个成熟系统中稳步推进高可用能力建设。【免费下载链接】bytebaseDatabase governance built for humans and agents — controlling changes and access across every major database.项目地址: https://gitcode.com/GitHub_Trending/by/bytebase创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表