版本跟踪
VersionTrackerStep 与其他步骤一样,是 SDK 中的常用步骤。批次成功处理后,VersionTrackerStep 会调用 save_processor_status() 的 trait 实现。
ProcessorStatusSaver
Section titled “ProcessorStatusSaver”ProcessorStatusSaver trait 要求实现具有以下签名的 save_processor_status 方法:
async fn save_processor_status( &self, last_success_batch: &TransactionContext<()>, ) -> Result<(), ProcessorError>;应在此方法中写入检查点。如果写入 Postgres,可以使用 SDK 的 Postgres 实现。使用 enum 可通过不同方式进行进度检查点记录。SDK 的 Postgres 实现使用简单的 processor_status 模型插入数据。
处理器成功将版本跟踪信息写入选定存储后,重启时需要从该存储中检索最新成功版本。此处提供了返回最新已保存处理版本的 get_starting_version() 方法示例。随后可按如下方式使用 starting_version: u64。若没有检查点,处理器将从链的开头开始。
let transaction_stream = TransactionStreamStep::new(TransactionStreamConfig { starting_version: Some(starting_version), ..self.config.transaction_stream_config.clone() }) .await?;SDK 没有提供保存回填进度的 ProcessorStatusSaver 实现。要启用回填进度保存,需要更新 IndexerProcessorConfig、ProcessorStatusSaver 和 get_starting_version()。没有这些变更时,很难同时运行处于最新交易版本的实时处理器和回填处理器。
在 IndexerProcessorConfig 上添加用于 BackfillConfig 的额外字段。在该实现中,BackfillConfig 是 enum ProcessorMode 的一部分,用于确定处理器运行的模式。回填模式下,处理器从不同版本开始,进度保存到单独的表中。
config.yaml 更新
Section titled “config.yaml 更新”在 YAML 文件的 server_config 中添加 backfill_config 部分来设置 backfill_alias。请参阅示例。
回填处理器状态表
Section titled “回填处理器状态表”为回填处理器状态使用单独的表,以避免写入冲突。该表(backfill_processor_status_table)使用 backfill_alias 而非 processor_name 作为主键,从而在并发运行头部处理器和回填处理器时避免与主 processor_status 表冲突。创建具有不同 backfill_alias 和交易版本范围的多个回填处理器,可加快回填速度。在此实现基础上扩展。该模型引入 BackfillStatus 新状态:InProgress 或 Complete,用于确定回填重启行为。
更新 ProcessorStatusSaver
Section titled “更新 ProcessorStatusSaver”扩展 ProcessorStatusSaver 实现以包含 Backfill 变体:从 BackfillConfig 提取 backfill_alias,并从 IndexerProcessorConfig.transaction_stream_config 提取 backfill_start_version 和 backfill_end_version,示例如此。更新相应写入查询,将数据写入新的 backfill_processor_status 表。
更新 get_starting_version
Section titled “更新 get_starting_version”当 IndexerProcessorConfig 中存在 BackfillConfig 字段时,在 get_starting_version 方法中添加语句,以查询 backfill_processor_status_table。