
后端任务调度工作流自动化微服务【免费下载链接】cadenceCadence is a distributed, scalable, durable, and highly available orchestration engine to execute asynchronous long-running business logic in a scalable and resilient way.项目地址https://gitcode.com/gh_mirrors/cad/cadence点击查看免费下载Scanner扫描器与 Fixer修复器是 Cadence 内部用于数据清理与一致性修复的后台子系统Scanner 周期性全量扫描数据库通过一系列“不变量检查”Invariant Check发现孤儿历史分支、废弃 TaskList、损坏执行记录等问题Fixer 则消费 Scanner 的输出对确认损坏的数据执行修复。本篇以service/worker/scanner/README.md为主干结合service/worker/scanner与common/reconciliation下的源码实现系统讲解其工作流结构、全部动态配置键、本地验证方法、blobstore 结果文件解读以及当前实现已知的边界与坑帮助你安全地评估、启用并调试这类数据修复能力。重要声明本文内容继承自仓库内 service/worker/scanner/README.md。该目录下绝大多数代码默认被动态配置禁用整体应视为 beta 级质量。启用前请务必理解你正在开启什么风险自负。任何官方推荐的修复流程只会在 Release Notes 等渠道明确标出本文档仅作描述不作推荐。此外文档作者明确指出Scanner/Fixer 相关代码存在诸多问题不应当视为项目希望长期保留的结构本文的价值在于把“踩坑经验”沉淀下来节省后人理解时间。这个目录是干什么的service/worker/scanner目录整体上承载了多种数据清理工作流data-cleanup workflows典型任务包括找出旧的无用数据并删除找出由历史 bug 造成的数据并修复清理并移除废弃的 tasklist避免其持续占用空间。作为一般规则这些工作流都会扫描整个数据库针对某一种数据检查某些条件必要时进行清理。但每种工作流“怎么做”差异很大。例如history scavenger历史清扫器负责找出在更新 workflow 官方历史时输掉 CAS 竞争而遗留的旧历史分支。由于任何清理都可能失败无法保证在工作流运行期间这些残留能被清掉因此 history scavenger 周期性遍历整个数据库查找这些孤儿历史分支并删除它们。其中最复杂的流程都建立在Scanner与Fixer之上因此这份 README 几乎只为它们而写。其余流程tasklist scanner、history scavenger、data corruption workflow 等行为更局部、更简单直接阅读代码比读任何文档都快例如 service/worker/scanner/workflow.go 和 service/worker/scanner/data_corruption_workflow.go。Scanner 与 Fixer 工作流的基本结构Invariant定义“什么是对的”以及如何修Invariant不变量定义Check和Fix两个方法Check用于确认不变量是否成立Fix用于在不成立时进行修复。它们位于 common/reconciliation/invariant 目录例如concreteExecutionExists.goCheck检查 current execution 记录是否指向一个真实存在的 concrete record如果不存在Fix会删除该 current recordhistory_exists.go检查执行记录是否具备应有的历史open_current_execution.go检查“当前执行”记录是否处于正确状态stale_workflow.go检查已完成的 workflow 是否超出保留期retention 10 天缓冲而仍未被清理此外还有history_invalid.go、inactive_domain_exists.go、mismatched_records.go、timer_invalid.go等。部分 Invariant 有配套的 “invariant collection”目前是 1:1 关系用于在数据类型并非唯一时例如 timer 只有一种用名称引用一组 Invariant。多个 Invariant 常常被打包进一个InvariantManager由它统一运行并聚合结果。其聚合规则可以在 invariant_manager.go 中看到RunChecks依次运行每个 Invariant 的Check结果的判定类型按healthy - corrupted - failed的优先级“劣化”聚合一旦出现 corrupted整体即 corrupted一旦出现 failed整体即 failednextCheckResultType的实现RunFixes依次运行每个 Invariant 的Fix聚合类型按skipped - fixed - failed的优先级推进nextFixResultType。Invariant 几乎只被 Scanner 和 Fixer 使用唯一例外是 common/ndc/history_resender.go 在 replication 处理中用到其中一个但那与 Scanner/Fixer 基本无关。Scanner只 Check把失败项推入 blobstoreScanner 只运行Check并把所有失败的检查结果推入 blobstore。要点如下核心数据来自一个Iterator其实现取决于你的 datastoreSQL / NoSQL 各自提供分页迭代器对每个 Iterator 条目通过InvariantManager运行一整套 Invariant每个条目的聚合检查结果被收集起来推入 blobstore。从源码看ShardScanner 持有一个failedWriter与corruptedWriter均为 blobstore writerScan方法遍历pagination.IteratorCheckResultTypeHealthy什么都不做CheckResultTypeCorrupted写入corruptedWriter并累计 corruption 统计CheckResultTypeFailed写入failedWriter并累计 check failed 统计。最终ScanReport包含每个 shard 的ScanKeys指向已 flush 的 corrupted/failed blob 文件或ControlFlowFailure遇到无法继续尝试检查/修复的严重错误。Fixer只 Fix消费最近一次 Scanner 的输出Fixer 只运行Fix作用于最近一次 Scanner 下载下来的结果。结构上与 Scanner 极其相似区别在于它调用Fix而非Check它使用另一个 Iterator专门迭代blobstore 中的 scanner 结果store.NewBlobstoreIterator见 executions/concrete_execution.go 的concreteExecutionFixerIterator所有已配置的 Invariant 都会运行而不只是上次Check失败的那些因为只调用FixInvariant 通常应该先在内部自行Check例如 stale_workflow 的 Fix 会先判断是否还在保留期内。源码中 ShardFixer 持有skippedWriter、failedWriter、fixedWriter三个 writer根据FixResultTypeFixed/Skipped/Failed分别落盘同时它通过allowDomain动态配置*FixerDomainAllow决定某个 domain 的数据是否允许被修复——不允许的 domain 一律记为 skipped。工作流如何被编排以上内容由*shardscanner.ScannerConfig实例装配起来它包含定制每一种 scanner/fixer 所需的全部信息并以 workflow 类型名作为 key 存入 / 取自 context。参见 executions/concrete_execution.go 的ConcreteExecutionConfig其中包含 workflow 类型名注册的函数名、启动参数、控制 scanner/fixer 工作流的高层动态配置enabled、concurrency 等以及 scanner 与 fixer 的 “hooks”工作流本身基本不关心这份配置它们每轮都执行相同的 activity由 activity 自己决定怎么做Activity 通过GetScannerContext/GetFixerContext见 shardscanner/types.go取得配置与其他依赖这两个函数会按当前 activity 所属 workflow 类型从后台 context 中取出对应的ScannerContext/FixerContext。HooksScannerConfig中的字段承载了大量非工作流行为。例如 concrete scanner 的 hooks 把“manager”InvariantManager、“iterator”遍历数据源、逐条产出待检查条目和“custom config”控制哪些 Invariant 启用的动态配置打包进concreteExecutionScannerHooks/concreteExecutionFixerHooks见 executions/concrete_execution.go。最终工作流在 activity 中根据配置创建 Scanner 或 Fixerscanner_workflow.go 中的ScannerWorkflow.Start先通过ActivityScannerConfig解析配置含动态配置覆盖项若未启用则直接返回随后按Concurrency并发地把 shard 列表分批ActivityBatchSize交给ActivityScanShardscanShardActivity后者逐 shard 调用scanShard由NewScanner依据参数/环境/hooks 创建真正的 Scanner 实例全部批次完成后执行ActivityScannerEmitMetrics上报指标Fixer 与其非常相似区别在于FixerWorkflow.Start会先查询上一次 Scanner 的运行结果拿到需要处理的 blobstore 文件列表GetCorruptedKeys再按Concurrency与ActivityBatchSize分批执行ActivityFixShard。Fixer 的默认并发等参数在 fixer_workflow.go 的resolveFixerConfig中定义Concurrency: 25、BlobstoreFlushThreshold: 1000、ActivityBatchSize: 200可被FixerWorkflowConfigOverwrites覆盖。Scanner/Fixer 的查询处理器同样注册在工作流上scanner_workflow.go 的getScanHandlers与 fixer_workflow.go 的setHandlers包括shard_report、shard_status、shard_status_summary、aggregate_report、domain_report、all_results等用于在 Web UI / CLI 中查看进度与统计。工作流的启动方式与生命周期运行 Scanner 和 Fixer 的工作流在服务启动时如果被启用即开始以每分钟一次、几乎永不过期的 cron形式持续运行每种 scanner/fixer 类型使用自己的 tasklist且只有启用时对应的 worker 才会被启动。这一点可以在 service/worker/scanner/scanner.go 的Start函数中看到遍历ShardScanners配置逐个判断ScannerEnabled()/FixerEnabled()仅为启用的配置注册 workflow 与 worker tasklist这意味着启用某项后服务重启即可立即生效但禁用只是暂停如果长时间后再次恢复可能并不符合预期因为是 cron 工作流只保留最初的启动参数后续对配置的修改不会自动生效如果你修改了相关字段请手动 cancel 或 terminate 对应 cron workflow然后重新启动 worker 来开启新版本。另外每个工作流只处理一种数据源主要通过其 Iterator 实现这意味着一个 scanner/fixer 内的所有 Invariant 处理的是同一类数据。例如 concrete executions scanner 只扫 concrete execution 记录见 executions/types.go 中的ConcreteExecutionType。配置如何启用 Scanner / FixerScanner 和 Fixer 默认是禁用的因为它们会消耗大量资源并可能为修正问题而删除或修改数据。因此通常需要修改动态配置才能运行它们。当前文档写作时可以通过如下配置启用。启用 scanner 工作流按数据源 / 记录类型如 concrete executions、timersworker.executionsScannerEnabled: - value: true # default false worker.currentExecutionsScannerEnabled: - value: true # default false worker.timersScannerEnabled: - value: true # default false worker.historyScannerEnabled: - value: true # default false worker.taskListScannerEnabled: - value: true # default true, only used on sql stores启用 scanner 的 Invariant目前每个 Invariant 只支持一种数据源 / 记录类型但同一数据源可以有多个 Invariant# concretes worker.executionsScannerInvariantCollectionStale: - value: true # default false worker.executionsScannerInvariantCollectionMutableState: - value: true # default true worker.executionsScannerInvariantCollectionHistory: - value: true # default true # timer invariant 是隐式的因为只有一种启用对应工作流即可。 # currents以下这些由于类型不匹配根本无法工作 worker.currentExecutionsScannerInvariantCollectionHistory: - value: true # default true worker.currentExecutionsInvariantCollectionMutableState: - value: true # default true这些*InvariantCollection*动态配置的解析逻辑在 executions/concrete_execution.go 的concreteExecutionCustomScannerConfig中每个为 true 的 collection 会以invariant.CollectionXXX.String()为 key 写入CustomScannerConfig随后由concreteExecutionScannerManager通过ParseCollectionsConcreteExecutionType.ToInvariants装配出实际运行的 Invariant 列表。启用 fixer 工作流同样每种类型一个worker.concreteExecutionFixerEnabled: - value: true # default false worker.currentExecutionFixerEnabled: - value: true # default false worker.timersFixerEnabled: - value: true # default false启用 fixer 在某个 domain 上运行必须配置否则任何数据都不会被修复worker.currentExecutionFixerDomainAllow: - value: true # default false constraints: {domainName: your-domain} # 例如或不带 constraints 以对所有 domain 生效 worker.concreteExecutionFixerDomainAllow: - value: true # default false worker.timersFixerDomainAllow: - value: true # default false注意*FixerDomainAllow这类键在源码中是通过dynamicproperties.BoolPropertyFnWithDomainFilter带 domain 过滤的动态布尔属性读取的对应 fixer.go 中ShardFixer.allowDomain的判断未在允许列表中的 domain其数据一律被跳过skipped不会执行任何修复。启用 fixer 的 Invariant# concretes worker.executionsFixerInvariantCollectionStale: - value: true # default false worker.executionsFixerInvariantCollectionMutableState: - value: true # default true worker.executionsFixerInvariantCollectionHistory: - value: true # default true # timer invariant 在启用 timer-fixer 时随之启用只有一种 # current execution fixer 从未正常工作过目前也不支持动态配置与 scanner 不同的是fixer 侧的自定义配置解析concreteExecutionCustomFixerConfig无论 true 还是 false 都会写入 key这是为了与更早版本服务器“无此配置”的行为区分开见 executions/concrete_execution.go 中的注释说明。与动态配置键对应的底层源码上述所有键对应的dynamicproperties定义可以在 common/dynamicconfig/dynamicproperties 目录中找到如ConcreteExecutionsScannerEnabled、ConcreteExecutionFixerEnabled、ConcreteExecutionsScannerInvariantCollectionHistory、ConcreteExecutionFixerDomainAllow、ConcreteExecutionsScannerConcurrency、ConcreteExecutionsScannerPersistencePageSize、ConcreteExecutionsScannerBlobstoreFlushThreshold、ConcreteExecutionsScannerActivityBatchSize等并在 executions/concrete_execution.go 的ConcreteExecutionConfig中被统一装配。concrete execution scanner/fixer 的默认 cron 周期是*/5 * * * *每 5 分钟ExecutionStartToCloseTimeout 为 20 年量级并允许重复启动WorkflowIDReusePolicyAllowDuplicate。本地验证搭一个可调试的环境文档作者在阅读和修改这部分代码时采用了一套本地验证流程步骤如下启动默认的 docker-compose 集群docker/docker-compose.yml在本地对 scanner / fixer 做配置、代码等修改执行make cadence-server确认可以构建运行./cadence-server start --services worker启动一个 worker连接到你刚起的 docker compose 集群通过 Web UI 浏览工作流历史通常是 http://localhost:8088/domains/cadence-system/workflows。默认的 docker compose 配置会启动一个 worker 实例但由于默认动态配置中除了worker.taskListScannerEnabled外其余全部禁用容器内的 worker不会运行大部分scanner/fixer 的 tasklist也就不会抢走本地 worker 的任务。因此你通常可以不做任何改动直接运行然后在 docker 之外启动本地定制过的worker 服务一切就能正常工作——这样既可以利用现成的 docker compose yaml又可以用 Web UI 观察结果还能快速改代码、重新构建、重新运行、调试完全不用碰 docker。如果你确实需要调试 tasklist scanner建议制作一个自定义构建并修改 docker compose 文件使用你的构建。细节见下文但对其他 scanner/fixer不是必需的。仅 tasklist scanner 需要的 docker compose 修改实现方式有几种文档作者偏好修改docker-compose*.yaml让它使用自定义本地构建并通过动态配置禁用 tasklist scanner。具体步骤参见 docker/README.md。作者个人偏好使用一个独立的 auto-setup 镜像 tag避免影响以后未经定制的 docker compose 运行。例如services: # ... cadence: image: ubercadence/server:my-auto-setup # 使用你自己的新 tag # ... environment: # 注意这个环境变量它就是你需要修改的文件 - DYNAMIC_CONFIG_FILE_PATHconfig/dynamicconfig/development.yaml修改动态配置文件后重新构建即可配置文件会被复制进 docker 镜像本地改动不会影响正在运行的容器。worker.taskListScannerEnabled: - value: false # default true, only used on sql stores只需设置上述配置并确保其他配置没有显式启用它们默认就是禁用的通常就完成了。docker 之外还需要改什么配置本地运行 server 一般使用config/dynamicconfig/development.yaml文件因此你很可能需要修改它来启用你的代码数据为了让你的 invariant / scanner 有东西可查运行一些 workflow然后在数据库中手动删除/修改数据是最简单的办法你的 Invariant为了防止过早地清除你手工改动的数据文档作者很推崇这个技巧——在Check里注入假失败让所有记录都流向 fixer而Fix里只打印而不真正修复func (i *invariant) Check(...) { x : true // go vet 目前不会对这种死代码报警。很方便 if x { return CheckResult{Failed} // 假失败让所有记录进入 fixer } // ... 其余正常代码保持不变 } func (i *invariant) Fix(...) { // 只打印修复动作而不真正执行这样下一次运行还会再尝试。 // 或者在这里也用 if x { 技巧 }IDE本地启动只带 worker 即可其他服务并非必需只会拖慢启动/关闭确保启动参数里有start --services worker。运行全部流程并检查结果加上断点或打印语句跑起来看看会发生什么。如果 Scanner 发现了有趣的东西你现在应该有一个/tmp/blobstore目录里面有{uuid()}_0.corrupted之类的文件这些 uuid 是随机的_0.corrupted后缀表示这是第 0 页共 N 页并且内容指向损坏条目每个发现问题的 scanner shard 对应一个 uuidshard 数可配置用于控制并发如果结果超过分页大小上限每个 shard 会有多个页。如果 Fixer 发现了什么/tmp/blobstore里现在应该有{uuid()}_0.skipped和/或{uuid()}_0.fixed文件这些 uuid同样是随机的并不指向它们数据来源的那个 Scanner 文件uuid 和分页模式与 Scanner 相同。你也许还会看到*.failed文件遵循同样的模式。这些是 Invariant 返回 Failure 结果的案例可能来自 scanner或fixer。不过只有 scanner 产生的*.corrupted文件会被 fixer 处理。注意*.failed文件可能包含所有状态的 invariant 结果因为一条记录的状态是“趋向于 failed”的只记录最终状态。具体行为参见 common/reconciliation/invariant/invariant_manager.goRunChecks/RunFixes的聚合逻辑。你也可以在 UI 中直接查看 scanner 与 fixer 工作流的结果。具体来说Scanner每种 scanner 类型有唯一 ID例如concreteExecutionsScannerWFID cadence-sys-executions-scanner查看 activities 可以看到每个 shard 发现了多少 corruptionQuery 它的aggregate_reportUI 可用其他查询需要参数目前需要用 CLI获取总体计数activities 返回的结果带有 UUID 和页码对应/tmp/blobstore中文件的 UUID 和页范围据此可以查到具体的详细结果否则作者的做法是拿一个已知的 workflow ID在最近一批文件里 grep 它。Fixer查看最近的 fixer 工作流获取修复结果如果某个 fixer 已经完成它很可能不是最近运行或正在运行的那个——往更早的看直到找到事件数超过十几个的那些才是 no-opactivities 接收来自 scanner 的 UUID 与页范围与 scanner 的返回值对应指向/tmp/blobstore中的扫描结果文件并返回同样结构的结果新的随机 UUID、新的页范围指向/tmp/blobstore中的新文件同样作者也推荐直接 grep 那些本应被处理的已知 ID。如果你没有打印或调试你关心的信息请检查这些文件的内容确认它们的行为符合预期。一个可工作的 scanner/fixer 配置示例体现在 /tmp/blobstore 文件中这是新的 stale invariant 在 concrete execution scanner - concrete execution fixer 中工作的示例其中用假结果简化测试作者当时也把其他 concrete invariants 都打开了出于好奇。首先做了这样的代码改动让Check永远失败而Fix运行真正的检查func (c *staleWorkflowCheck) Check( ctx context.Context, execution interface{}, ) CheckResult { x : true if x { return CheckResult{ CheckResultType: CheckResultTypeCorrupted, InvariantName: c.Name(), Info: fake corrupt, } } _, result : c.check(ctx, execution) return result }并加入如下动态配置# 启用这些 invariants worker.executionsScannerInvariantCollectionStale: - value: true # default false worker.executionsScannerInvariantCollectionMutableState: - value: true # default true worker.executionsScannerInvariantCollectionHistory: - value: true # default true worker.executionsFixerInvariantCollectionStale: - value: true # default false worker.executionsFixerInvariantCollectionMutableState: - value: true # default true worker.executionsFixerInvariantCollectionHistory: - value: true # default true # 这些 invariants 都由 concrete execution scanner/fixer 运行 worker.executionsScannerEnabled: # 注意名称略有不同 - value: true # default false worker.concreteExecutionFixerEnabled: - value: true # default false worker.concreteExecutionFixerDomainAllow: - value: true # default false一次 scanner 与 fixer 运行之后/tmp/blobstore里出现了*.corrupted和*.skipped文件。*.corrupted文件内容类似{ Input: { Execution: { ... }, Result: { CheckResultType: corrupted, DeterminingInvariantType: stale_workflow, CheckResults: [ { CheckResultType: healthy, InvariantName: history_exists, Info: , InfoDetails: }, { CheckResultType: healthy, InvariantName: open_current_execution, Info: , InfoDetails: }, { CheckResultType: corrupted, InvariantName: stale_workflow, Info: fake corrupt, InfoDetails: } ] } }, Result: { FixResultType: skipped, DeterminingInvariantName: null, FixResults: null } }可以看到两个 healthy invariants以及一个被伪造的 corrupted。随后 fixer 运行时在*.skipped文件中得到{ Execution: { ... }, Input: { Execution: { ... }, Result: { CheckResultType: corrupted, DeterminingInvariantType: stale_workflow, CheckResults: [ { CheckResultType: healthy, InvariantName: history_exists, Info: , InfoDetails: }, { CheckResultType: healthy, InvariantName: open_current_execution, Info: , InfoDetails: }, { CheckResultType: corrupted, InvariantName: stale_workflow, Info: fake corrupt, InfoDetails: } ] } }, Result: { FixResultType: skipped, DeterminingInvariantName: null, FixResults: [ { FixResultType: skipped, InvariantName: history_exists, CheckResult: { CheckResultType: healthy, InvariantName: history_exists, Info: , InfoDetails: }, Info: skipped fix because execution was healthy, InfoDetails: }, { FixResultType: skipped, InvariantName: open_current_execution, CheckResult: { CheckResultType: healthy, InvariantName: open_current_execution, Info: , InfoDetails: }, Info: skipped fix because execution was healthy, InfoDetails: }, { FixResultType: skipped, InvariantName: stale_workflow, CheckResult: { CheckResultType: , InvariantName: , Info: , InfoDetails: }, Info: no need to fix: completed workflow still within retention 10-day buffer, InfoDetails: completed workflow still within retention 10-day buffer, closed 2023-09-20 20:26:04.924876012 -0500 CDT and allowed to exist until 2023-10-07 } ] } }注意fixer 中全部三个 invariants 都运行了三个都启用了且因为没发现任何问题三个修复都被跳过。如果当时也伪造了 stale workflow invariant 的Fix你会在 fixer 中看到该 invariant 的FixResultType为 fixed文件名也会是*.fixed而不是*.skipped。一个错误配置的反例类型不匹配这是 current-execution scanner 工作、而 current-execution fixer 行为异常并使用错误类型、产生*.failed文件的案例对应 commiteb55629d时的状态。作者伪造了 current execution invariantCheck永远返回 corruptFix直接 panic并在动态配置中启用 current execution scanner 和 fixer然后运行 worker。首先scanner 运行产生*.corrupted文件条目如下{ Input: { Execution: { ... }, Result: { CheckResultType: corrupted, DeterminingInvariantType: concrete_execution_exists, CheckResults: [ { CheckResultType: corrupted, InvariantName: concrete_execution_exists, Info: execution is open without having concrete execution, InfoDetails: concrete execution not found. WorkflowId: e905c98f-108a-4191-9ef2-ca07a1361f9c, RunId: 6bc5386b-c043-4eb1-a332-c3bb7b5188f0 } ] } }, Result: { FixResultType: skipped, DeterminingInvariantName: null, FixResults: null } }这证明 scanner 正确识别了 always corrupt 的结果对应 common/reconciliation/invariant/concreteExecutionExists.go。这些数据随后被 current execution fixer 消费产出*.failed文件内容如下{ Execution: { ... }, Input: { Execution: { ... }, Result: { CheckResultType: corrupted, DeterminingInvariantType: concrete_execution_exists, CheckResults: [ { CheckResultType: corrupted, InvariantName: concrete_execution_exists, Info: execution is open without having concrete execution, InfoDetails: concrete execution not found. WorkflowId: e905c98f-108a-4191-9ef2-ca07a1361f9c, RunId: 6bc5386b-c043-4eb1-a332-c3bb7b5188f0 } ] } }, Result: { FixResultType: failed, DeterminingInvariantName: history_exists, FixResults: [ { FixResultType: failed, InvariantName: history_exists, CheckResult: { CheckResultType: failed, InvariantName: history_exists, Info: failed to check: expected concrete execution, InfoDetails: }, Info: failed fix because check failed, InfoDetails: }, { FixResultType: failed, InvariantName: open_current_execution, CheckResult: { CheckResultType: failed, InvariantName: open_current_execution, Info: failed to check: expected concrete execution, InfoDetails: }, Info: failed fix because check failed, InfoDetails: } ] } }可以清楚看到scanner 的数据作为输入结果被带入而fixer 侧运行并失败的却是完全不同的 invariantshistory_exists、open_current_execution。这正是文档与 executions/current_execution.go 顶部注释所警告的问题current execution fixer 被写成使用 concrete execution 的 invariants导致它一旦运行就必然失败。已知问题与注意事项最后文档作者对这套实现的总结性观察“Thoughts and prayers”值得每一位准备启用的人认真阅读外部扩展极其困难当前实现重度依赖无法新增或修改的常量修改它很可能需要重写大量代码。未来版本应当修复这一点——自定义数据库插件可能有独特的问题需要独特的 scan/fix 工具。current-execution fixer 从未在任何地方成功运行过scanner 看起来可用、invariant 看起来正确但 fixer 使用了 concrete execution 的 invariants导致运行必败。不要从它的代码或配置得出任何结论。它可能在未来被修复并可选地启用或被删除。current / concrete / timer 三套 scan/fix 代码非常相似但不完全相同可能好也可能不好但确实容易让人困惑阅读时要仔细。配置键与可配置性在各实现之间差异很大复制粘贴不要猜测并且要验证之后再宣称“完成”。先扫全部、再修全部 的结构在大集群上慢得成问题且由于所有数据至少处理 2 遍而相当浪费资源。放慢速度重新检查本身未必是坏事但大集群上 fixer 可能要数周之后才开始这对快速验证修复或及时处理问题非常糟糕。用 Query 在 scanner 与 fixer 之间传递数据显得很别扭作者怀疑这是为了避免返回过多数据而违反 blob 大小限制。但由于实际使用的数据是 fixer 的查询 activity 保存的它仍然受该限制约束。而且结果分散在大量 activities 与 queries 之间很难获得全局概览——或许可以把全部结果推入一个新的 blobstore 文件供人类查看。整体存在大量难以解释的间接跳转导致控制流极难梳理可能是为了尚未用到的灵活性可能只是需要现代化改造用实例字段替代后台 activity context也可能是当年 Go 缺少泛型时的一种函数式风格的无奈之举。总体而言这套代码可能值得重写不过其中一些部件invariants、iterators 等明显可以复用。进一步阅读service/worker/scanner/README.md本篇文章的原始出处含大量第一手踩坑记录service/worker/scanner/shardscanner/scanner.go 与 service/worker/scanner/shardscanner/fixer.goScanner / Fixer 核心实现service/worker/scanner/shardscanner/scanner_workflow.go 与 service/worker/scanner/shardscanner/fixer_workflow.go工作流编排、并发与查询处理器service/worker/scanner/shardscanner/types.goScanner/Fixer 的参数、报告与结果结构service/worker/scanner/executions/concrete_execution.go 与 service/worker/scanner/executions/current_execution.go具体类型的配置装配注意后者顶部对 current execution fixer 的警告common/reconciliation/invariant全部 Invariant 实现与 invariant_manager.go 的聚合逻辑service/worker/scanner/scanner.goScanner 子系统的启动入口Start函数service/worker/scanner/workflow.go 与 service/worker/scanner/data_corruption_workflow.go不遵循上述模式的更简单工作流。赞分享后端任务调度工作流自动化微服务【免费下载链接】cadenceCadence is a distributed, scalable, durable, and highly available orchestration engine to execute asynchronous long-running business logic in a scalable and resilient way.项目地址https://gitcode.com/gh_mirrors/cad/cadence点击查看免费下载相关推荐用 updatecli 驱动 Rancher 依赖版本自动化定时工作流、Manifest 结构与本地验证实践用 updatecli 驱动 Rancher 依赖版本自动化定时工作流、Manifest 结构与本地验证实践 Rancher 作为一套完整的容器管理平台其仓云原生容器编排集群管理后端ADK 事件模型Event 与 NodeInfo深度指南会话、动作与工作流路由的底层数据结构ADK 事件模型Event 与 NodeInfo深度指南会话、动作与工作流路由的底层数据结构 在 Google Agent Development Kit人工智能AI AgentAgent 框架多智能体工具调用RAGagno AgentOS Cookbook 全量实测工作流01–24 课 Live 验证清单与结果契约深度解析agno AgentOS Cookbook 全量实测工作流01–24 课 Live 验证清单与结果契约深度解析 导读 本文以 cookbook/05_agen人工智能大模型AI AgentAgent 框架多智能体工具调用RAGAgent 工作流Agent 记忆上一篇Qwen-Image-Lightning模型评测轻量级AI绘图工具性能对比下一篇Amazing-QR代码解读combine函数图片合成原理详解创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考