跳转到内容

交易流步骤

TransactionStreamStep 是交易处理流水线的基础组件。它与 TransactionStream 服务建立 gRPC 连接,批量获取交易并将其输出以供后续处理。该步骤还会在暂时性故障发生时管理连接重试和重新连接。它通常是处理器的初始步骤,负责向下游步骤流式传送交易。

  1. 获取交易:从 gRPC 服务检索交易批次。
  2. 管理连接:处理 gRPC 重新连接,确保流具有弹性。
  3. 提供元数据:将版本和时间戳等上下文信息附加到交易。

TransactionStreamStep 结构体定义如下:

pub struct TransactionStreamStep
where
Self: Sized + Send + 'static,
{
transaction_stream_config: TransactionStreamConfig,
pub transaction_stream: Mutex<TransactionStreamInternal>,
}
  • TransactionStreamStep 连接到 gRPC TransactionStream 服务。
  • 它使用 poll 方法持续轮询新交易。
  • 每个批次封装在 TransactionContext 中,其中包括以下元数据:
    • 起始和结束版本。
    • 交易时间戳。
    • 以字节为单位的批次大小。
  • 若连接中断,它会尝试无缝重新连接。