跳转到内容

创建步骤

步骤是 SDK 中的处理逻辑单元,可用于定义数据提取、转换或存储逻辑。 步骤是处理器的构建模块。 Aptos core 处理器将以下操作分别表示为独立步骤:(1) 从交易流获取交易流,(2) 提取数据,(3) 写入数据库,以及 (4) 跟踪进度。

SDK 中有两类步骤:

  1. AsyncStep:处理一批输入项并返回一批输出项。
  2. PollableAsyncStep:与 AsyncStep 相同,但还会定期轮询其内部状态,并在有可用数据时返回一批输出项。

要使用 SDK 创建步骤,请按以下说明操作:

  1. 实现 Processable trait。该 trait 定义了步骤的几个重要细节:输入和输出类型、处理逻辑,以及运行类型(AsyncStepRunTypePollableAsyncStepRunType)。

    #[async_trait]
    impl Processable for MyExtractorStep {
    // The Input is a vec of Transaction
    type Input = Vec<Transaction>;
    // The Output is a vec of MyData
    type Output = Vec<MyData>;
    // Depending on the type of step this is, the RunType is either
    // - AsyncRunType
    // - PollableAsyncRunType
    type RunType = AsyncRunType;
    // Processes a batch of input items and returns a batch of output items.
    async fn process(
    &mut self,
    input: TransactionContext<Vec<Transaction>>,
    ) -> Result<Option<TransactionContext<Vec<MyData>>>, ProcessorError> {
    let transactions = input.data;
    let data = transactions.iter().map(|transaction| {
    // Define the processing logic to extract MyData from a Transaction
    }).collect();
    Ok(Some(TransactionContext {
    data,
    metadata: input.metadata,
    }))
    }
    }

    在上述示例中,输入和输出类型被封装在 TransactionContext 中。 TransactionContext 包含正在处理的数据批次的相关元数据,例如交易版本和时间戳,用于指标和日志记录。

  2. 实现 NamedStep trait。它用于日志记录。

    impl NamedStep for MyExtractorStep {
    fn name(&self) -> String {
    "MyExtractorStep".to_string()
    }
    }
  3. 实现 AsyncStep trait 或 PollableAsyncStep trait,以定义该步骤如何在处理器中运行。

    1. 如果使用 AsyncStep,请将以下代码加入代码库:

      impl AsyncStep for MyExtractorStep {}
    2. 如果创建 PollableAsyncStep,需要定义轮询间隔以及每次轮询时步骤应执行的操作。

      #[async_trait]
      impl<T: Send + 'static> PollableAsyncStep for MyPollStep<T>
      where
      Self: Sized + Send + Sync + 'static,
      T: Send + 'static,
      {
      fn poll_interval(&self) -> std::time::Duration {
      // Define duration
      }
      async fn poll(&mut self) -> Result<Option<Vec<TransactionContext<T>>>, ProcessorError> {
      // Define code here on what this step should do every time it polls
      // Optionally return a batch of output items
      }
      }

构建提取器步骤时,需要定义希望如何从交易中解析数据。 请在此处详细了解如何从交易中解析数据。

SDK 随附一组可用于构建处理器的常见步骤

  1. TransactionStreamStep 向处理器提供 Aptos 交易流。请在此处了解更多。
  2. TimedBufferStep 缓冲一批项,并定期轮询以将这些项释放给下一步骤。
  3. VersionTrackerStep 跟踪处理器进度并为其进度创建检查点。请在此处了解更多。
  4. OrderByVersionStep 按起始版本对交易上下文排序。它缓冲这些已排序上下文,并在每个轮询间隔释放它们。
  5. WriteRateLimitStep 限制每秒写入数据库的字节数。