交易流步骤
TransactionStreamStep 是交易处理流水线的基础组件。它与 TransactionStream 服务建立 gRPC 连接,批量获取交易并将其输出以供后续处理。该步骤还会在暂时性故障发生时管理连接重试和重新连接。它通常是处理器的初始步骤,负责向下游步骤流式传送交易。
- 获取交易:从 gRPC 服务检索交易批次。
- 管理连接:处理 gRPC 重新连接,确保流具有弹性。
- 提供元数据:将版本和时间戳等上下文信息附加到交易。
TransactionStreamStep 结构体定义如下:
pub struct TransactionStreamStepwhere Self: Sized + Send + 'static,{ transaction_stream_config: TransactionStreamConfig, pub transaction_stream: Mutex<TransactionStreamInternal>,}TransactionStreamStep连接到 gRPCTransactionStream服务。- 它使用
poll方法持续轮询新交易。 - 每个批次封装在
TransactionContext中,其中包括以下元数据:- 起始和结束版本。
- 交易时间戳。
- 以字节为单位的批次大小。
- 若连接中断,它会尝试无缝重新连接。