跳转到内容

版本跟踪

VersionTrackerStep 与其他步骤一样,是 SDK 中的常用步骤。批次成功处理后,VersionTrackerStep调用 save_processor_status() 的 trait 实现。

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 实现。要启用回填进度保存,需要更新 IndexerProcessorConfigProcessorStatusSaverget_starting_version()。没有这些变更时,很难同时运行处于最新交易版本的实时处理器和回填处理器。

IndexerProcessorConfig 上添加用于 BackfillConfig额外字段。在该实现中,BackfillConfig 是 enum ProcessorMode 的一部分,用于确定处理器运行的模式。回填模式下,处理器从不同版本开始,进度保存到单独的表中。

在 YAML 文件的 server_config 中添加 backfill_config 部分来设置 backfill_alias。请参阅示例

为回填处理器状态使用单独的表,以避免写入冲突。该表(backfill_processor_status_table)使用 backfill_alias 而非 processor_name 作为主键,从而在并发运行头部处理器和回填处理器时避免与主 processor_status 表冲突。创建具有不同 backfill_alias 和交易版本范围的多个回填处理器,可加快回填速度。在此实现基础上扩展。该模型引入 BackfillStatus 新状态:InProgressComplete,用于确定回填重启行为。

扩展 ProcessorStatusSaver 实现以包含 Backfill 变体:从 BackfillConfig 提取 backfill_alias,并从 IndexerProcessorConfig.transaction_stream_config 提取 backfill_start_versionbackfill_end_version示例如此。更新相应写入查询,将数据写入新的 backfill_processor_status 表。

IndexerProcessorConfig 中存在 BackfillConfig 字段时,在 get_starting_version 方法中添加语句,以查询 backfill_processor_status_table