链下索引器高可用部署指南 — MSG Chain
数据来源:MSG Chain 代码库核实
主网状态: No-Go — 当前 MSGChain 主网裁决为 No-Go,以下内容反映代码实际状态,不代表生产可用。
目录
- 概述
- 链下索引器架构
- MSG Chain 索引器拓扑
- 高可用设计
- 数据处理管道
- 存储选型
- 索引器故障转移
- 数据一致性
- 性能基准
- 监控与告警
- 部署实践
- 与 MSG Chain Indexer 数据平面的集成
- 总结
1. 概述
1.1 什么是链下索引器
链下索引器(Off-Chain Indexer)是一类将区块链原始数据拉取到链下数据库中进行结构化存储、高效查询和实时推送的服务。与链上直接查询相比,链下索引器解决了三个核心矛盾:
| 问题 | 链上查询 | 链下索引器 |
|---|---|---|
| 查询能力 | 仅支持键值查询 | 支持 SQL、GraphQL、全文搜索、聚合统计 |
| 性能 | 受限于共识层,TPS 有限 | 使用传统数据库,毫秒级响应 |
| 历史数据 | 节点修剪后丢失 | 长期保存,支持任意时间点回溯 |
在 MSG Chain 的 CosmWasm 生态中,合约事件通过 Tendermint WebSocket 或 REST API 暴露。链下索引器订阅这些事件流,解析并转换为业务实体(Agent、PaymentSession、DIDDocument 等),最终通过 GraphQL 或 REST API 对外暴露。
1.2 高可用的必要性
索引器一旦宕机,将导致:
- 前端/DApp 查询中断:用户无法查看 Agent 列表、支付状态、DID 信息
- 数据断层:未处理的区块形成 gap,恢复时需要重新同步
- 链下服务连锁故障:依赖索引器的 Agent API、通知服务、数据分析平台全部受影响
因此,开发参考级索引器必须满足以下 SLA:
| 指标 | 目标值 |
|---|---|
| 可用性 | 99.9%+(全年宕机 < 8.76 小时) |
| 数据延迟 | < 3 秒(从区块确认到索引完成) |
| RTO(恢复时间目标) | < 5 分钟 |
| RPO(恢复点目标) | < 1 个区块高度 |
1.3 本文阅读对象
- MSG Chain 基础设施运维工程师
- DApp 后端开发人员
- 索引器插件开发者
- 希望对索引管道进行性能调优的技术决策者
2. 链下索引器架构
2.1 逻辑分层
链下索引器在逻辑上分为四个层次:
+------------------------------------------------------------------+
| API 层 (GraphQL / REST) |
| GraphQL Yoga / Apollo Server / Express.js |
| DataLoader / N+1 消除 / 查询缓存 |
+------------------------------------------------------------------+
| 存储后端层 |
| PostgreSQL 15+ (主存储) | Redis 7+ (缓存) | S3 (冷归档) |
| TimescaleDB (时序) | | |
+------------------------------------------------------------------+
| 处理管道层 |
| Block Sync → Event Filter → State Derivation → Reorg Handler |
| Checkpoint Manager → Batch Writer |
+------------------------------------------------------------------+
| 事件来源层 |
| Tendermint RPC (WebSocket) | gRPC | Polling HTTP |
| msg-chain-1 节点集群 |
+------------------------------------------------------------------+
2.1.1 事件来源层
负责从 MSG Chain 节点获取原始区块和事件数据。支持三种接入模式:
| 模式 | 延迟 | 吞吐 | 可靠性 | 推荐场景 |
|---|---|---|---|---|
| WebSocket 订阅 | 实时(~500ms) | 中 | 弱(断连需重连) | 实时事件流 |
| REST Polling | 轮询间隔(1-5s) | 高 | 强(无状态) | 历史同步 |
| gRPC 流 | 实时(~200ms) | 高 | 中 | 高性能场景 |
2.1.2 处理管道层
核心数据处理引擎,包含以下组件:
- Block Sync:管理同步进度,记录当前已处理的高度
- Event Filter:根据合约地址和事件类型过滤出业务相关事件
- State Derivation:从原始事件推导出业务状态增量
- Reorg Handler:检测并处理链重组事件
- Checkpoint Manager:定期记录处理位置,支持故障恢复
- Batch Writer:批量写入存储后端,优化写入吞吐
2.1.3 存储后端层
- PostgreSQL:主存储,存储所有业务实体
- Redis:缓存热点查询、去重布隆过滤器、实时推送消息队列
- TimescaleDB:存储时序数据(区块时间、事件频率、Gas 统计)
- S3 兼容对象存储:冷数据归档、历史快照备份
2.1.4 API 层
对外提供查询接口,支持:
- GraphQL:类型安全、关联导航、实时订阅
- REST:传统 RESTful 端点,向后兼容
- WebSocket:实时事件推送
2.2 数据流全景
MSG Chain 节点 (msg-chain-1)
│
│ WebSocket / REST / gRPC
▼
┌──────────────────────────┐
│ Block Source │ ← 事件来源层
│ (TendermintClient) │
└────────┬─────────────────┘
│ Block + Events
▼
┌──────────────────────────┐
│ Event Normalizer │ ← 处理管道层
│ ┌──────────────────────┐ │
│ │ Block Fetcher │ │
│ │ Event Parser │ │
│ │ Event Filter │ │
│ │ State Derivation │ │
│ │ Reorg Detector │ │
│ │ Checkpoint Manager │ │
│ │ Batch Committer │ │
│ └──────────────────────┘ │
└────────┬─────────────────┘
│ Normalized Entities
▼
┌──────────────────────────┐
│ Storage Backend │ ← 存储后端层
│ ┌──────────────────────┐ │
│ │ PostgreSQL (主) │ │
│ │ Redis (缓存/队列) │ │
│ │ S3 (冷归档) │ │
│ │ TimescaleDB (时序) │ │
│ └──────────────────────┘ │
└────────┬─────────────────┘
│ Data
▼
┌──────────────────────────┐
│ API Server │ ← API 层
│ GraphQL / REST / WS │
│ DataLoader / Cache │
└────────┬─────────────────┘
│ Responses
▼
Frontend / DApp / SDK
2.3 关键设计决策
2.3.1 推模式 vs 拉模式
| 维度 | 推模式 (WebSocket) | 拉模式 (Polling) |
|---|---|---|
| 实现复杂度 | 高(需处理重连、背压) | 低 |
| 实时性 | 高 | 取决于轮询间隔 |
| 恢复机制 | 复杂(需重连后 catchup) | 简单(从断点继续) |
| 资源消耗 | 维持长连接 | 定期 HTTP 请求 |
建议:生产环境使用混合模式——WebSocket 接收实时事件,REST Polling 作为保底同步手段。
2.3.2 单管道 vs 多管道
对于 MSG Chain 的五类合约(Agent Registry、Agent Payment、DID Registry、Constitution、Micropayment),有两种索引策略:
| 策略 | 说明 | 适用场景 |
|---|---|---|
| 单管道全量索引 | 一个索引器实例处理所有合约事件 | 开发环境、小型部署 |
| 多管道分片索引 | 每个索引器只处理一个或多个合约 | 生产环境、高吞吐场景 |
多管道模式允许对不同合约的索引速度进行独立控制,且隔离故障影响范围。
3. MSG Chain 索引器拓扑
3.1 最小部署拓扑
┌──────────────────────┐
│ MSG Chain RPC │
│ msg-chain-1 │
│ ws://node:26657 │
│ http://node:1317 │
└──────┬───────────────┘
│
▼
┌──────────────────────┐
│ Indexer (单实例) │
│ - Block Sync │
│ - Event Process │
│ - State Derive │
└──────┬───────────────┘
│
▼
┌──────────────────────┐
│ PostgreSQL 15 │
│ (单机) │
└──────┬───────────────┘
│
▼
┌──────────────────────┐
│ GraphQL API │
│ (单实例) │
└──────────────────────┘
最小拓扑适合开发/测试环境,存在单点故障风险。
3.2 生产拓扑(高可用)
┌─────────────────────┐
│ Load Balancer │
│ (HAProxy / Nginx) │
│ msgchain.org │
└──────┬──────────┬───┘
│ │
┌────────────┘ └────────────┐
▼ ▼
┌──────────────────┐ ┌──────────────────┐
│ GraphQL API A │ │ GraphQL API B │
│ (active) │ │ (active) │
└──────┬───────────┘ └──────┬───────────┘
│ │
└──────────────┬──────────────────────┘
│
┌───────────┴───────────┐
│ │
▼ ▼
┌──────────────────┐ ┌──────────────────┐
│ PostgreSQL │ │ PostgreSQL │
│ Primary (读写) │──│ Replica (只读) │
│ - blocks │ │ - blocks │
│ - transactions │ │ - transactions │
│ - events │ │ - events │
│ - agents │ │ - agents │
│ - sessions │ │ - sessions │
│ - did_documents │ │ - did_documents │
│ - constitution │ │ - constitution │
│ - channels │ │ - channels │
└──────────────────┘ └──────────────────┘
│
│
┌───────────┴───────────┐
│ │
▼ ▼
┌─────────────────────┐ ┌─────────────────────┐
│ Indexer Primary │ │ Indexer Standby │
│ (active - leader) │ │ (passive - follow) │
│ │ │ │
│ - 连接 RPC │ │ - 健康监控 │
│ - 处理 Event │ │ - 同步 Standby DB │
│ - 写入 Primary DB │ │ - 冷备就绪 │
│ - 维护 Checkpoint │ │ │
└──────┬──────────────┘ └─────────────────────┘
│ │
└────────────┬───────────┘
│
▼
┌──────────────────────┐
│ Redis 集群 │
│ - 缓存 │
│ - 布隆过滤器 │
│ - 消息发布/订阅 │
│ - Leader 选举锁 │
└──────────────────────┘
│
▼
┌──────────────────────┐
│ S3 兼容对象存储 │
│ (冷归档 / 备份) │
└──────────────────────┘
3.3 组件说明
| 组件 | 规格建议 | 角色 |
|---|---|---|
| Indexer Primary | 4C8G × 1 | 主索引器,执行数据处理管道 |
| Indexer Standby | 4C8G × 1 | 备用索引器,接管时启动处理 |
| GraphQL API | 2C4G × 2 | 无状态查询服务,水平扩展 |
| PostgreSQL Primary | 8C32G, 500GB SSD | 主数据库,写入业务数据 |
| PostgreSQL Replica | 8C32G, 500GB SSD | 只读副本,分担查询压力 |
| Redis | 4C8G × 3 (集群) | 缓存、锁、消息队列 |
| Load Balancer | 2C4G × 2 | 流量分发 + TLS 终结 |
3.4 网络拓扑
Internet
│
▼
┌──────────────────────────────────────────┐
│ DMZ 网络 │
│ msgchain.org → HAProxy (80/443) │
│ → GraphQL API A / B │
└──────────────────────────────────────────┘
│
▼
┌──────────────────────────────────────────┐
│ 内部网络 │
│ Indexer Primary ↔ PostgreSQL Primary │
│ Indexer Standby ↔ PostgreSQL Replica │
│ Indexer → MSG Chain RPC (内网节点) │
│ All → Redis Cluster │
│ All → S3 (对象存储) │
└──────────────────────────────────────────┘
4. 高可用设计
4.1 高可用模式对比
索引器高可用有三种典型模式:
| 模式 | 切换延迟 | 数据一致性 | 资源消耗 | 复杂度 |
|---|---|---|---|---|
| Active-Passive | 10-30s | 强一致 | 2N | 低 |
| Active-Active | 0s(故障自动分发) | 最终一致 | 2N+ | 高 |
| 分片索引 | 0s(单分片故障不影响其他) | 分片内强一致 | N+ 副本 | 中 |
4.1.1 Active-Passive(推荐)
正常时:
Primary (Leader) → 处理事件 → 写入 Primary DB → 更新 Checkpoint
Standby (Follower) → 周期性心跳 → 监控 Checkpoint → 拉取 Standby DB
故障时:
1. Standby 检测到 Primary 心跳超时 (> 3 个周期)
2. Standby 尝试获取 Redis 分布式锁
3. 获取成功后: Standby 晋升为 Primary
4. 从最后一个 Checkpoint 位置继续处理
5. 重新健康检查 → 通知负载均衡器切换
核心机制:
// 伪代码:Leader 选举协议
interface ElectionConfig {
lockKey: "indexer:leader:lock";
lockTTL: 15; // 锁持有时间(秒)
heartbeatInterval: 5; // 心跳间隔(秒)
retryInterval: 3; // 选举重试间隔(秒)
maxRetries: 5;
}
async function tryAcquireLeadership(): boolean {
const locked = await redis.setNX(lockKey, instanceId, lockTTL);
if (!locked) return false;
// 启动心跳续约协程
startHeartbeat(lockKey, instanceId, lockTTL, heartbeatInterval);
return true;
}
async function extendLeadership(): void {
while (isLeader) {
await redis.expire(lockKey, lockTTL);
await sleep(heartbeatInterval * 1000);
}
}
4.1.2 Active-Active
┌──────────────┐ ┌──────────────┐
│ Indexer A │ │ Indexer B │
│ (active) │ │ (active) │
│ │ │ │
│ height 1-50 │ │ height 1-50 │
│ height │ │ height │
│ 101-150 │ │ 51-100 │
└──────┬───────┘ └──────┬───────┘
│ │
└──────────┬─────────────┘
│
▼
┌────────────────┐
│ PostgreSQL │
│ (写入冲突?) │
└────────────────┘
Active-Active 需要解决写冲突问题:
- 按合约分片:Indexer A 处理 Agent Registry + Agent Payment,Indexer B 处理 DID + Constitution + Micropayment
- 按高度水平分片:Indexer A 处理偶数高度,Indexer B 处理奇数高度(不推荐,跨块查询复杂)
- 幂等写入:同一事件被两个实例处理时,使用 UPSERT + 乐观锁去重
由于 CosmWasm 事件处理的天然有序性要求(后一个事件依赖前一个事件的状态),MSG Chain 推荐 Active-Passive 模式,辅以分片策略处理吞吐瓶颈。
4.1.3 分片索引
┌─────────────────────────────────────────────┐
│ Indexer Orchestrator │
│ (分配分片 / 监控健康 / 动态重平衡) │
└────┬────────┬────────┬────────┬─────────────┘
│ │ │ │
▼ ▼ ▼ ▼
┌────────┐┌────────┐┌────────┐┌────────────┐
│Shard 0 ││Shard 1 ││Shard 2 ││Shard 3 │
│Agent ││Payment ││DID ││Constitution│
│Registry││Session ││Registry││+ Micropay │
└───┬────┘└───┬────┘└───┬────┘└──────┬─────┘
│ │ │ │
└─────────┴─────────┴────────────┘
│
▼
┌──────────────────┐
│ PostgreSQL │
│ (共享数据库) │
│ (按 shard_id 分区)│
└──────────────────┘
4.2 Leader 选举详细实现
使用 Redis 分布式锁 + TTL 续约 实现 Leader 选举:
[实例启动]
│
▼
[尝试获取 Redis 锁: SET indexer:leader:{chainId} {instanceId} NX EX 15]
│
├── 成功 ──→ [成为 Leader] ──→ [启动处理协程] ──→ [每 5s 续约锁]
│
└── 失败 ──→ [成为 Follower] ──→ [启动监控协程]
│
▼
[每 3s 检查 Leader 锁]
│
Leader 锁过期?
│
├── 是 ──→ [尝试成为 Leader]
│
└── 否 ──→ [继续等待]
// src/ha/leader-election.ts
import Redis from 'ioredis';
interface LeaderElectionOptions {
lockKey: string;
instanceId: string;
lockTTL: number;
heartbeatInterval: number;
onElected: () => Promise<void>;
onRevoked: () => Promise<void>;
}
export class LeaderElection {
private redis: Redis;
private isLeader = false;
private heartbeatTimer: NodeJS.Timeout | null = null;
private monitorTimer: NodeJS.Timeout | null = null;
private options: LeaderElectionOptions;
constructor(redisUrl: string, options: LeaderElectionOptions) {
this.redis = new Redis(redisUrl);
this.options = options;
}
async start(): Promise<void> {
await this.tryAcquireLeadership();
}
private async tryAcquireLeadership(): Promise<void> {
const acquired = await this.acquireLock();
if (acquired) {
this.isLeader = true;
this.startHeartbeat();
await this.options.onElected();
} else {
this.startMonitoring();
}
}
private async acquireLock(): Promise<boolean> {
const result = await this.redis.set(
this.options.lockKey,
this.options.instanceId,
'NX',
'EX',
this.options.lockTTL,
);
return result === 'OK';
}
private async releaseLock(): Promise<void> {
const script = `
if redis.call("GET", KEYS[1]) == ARGV[1] then
return redis.call("DEL", KEYS[1])
end
return 0
`;
await this.redis.eval(script, 1, this.options.lockKey, this.options.instanceId);
}
private startHeartbeat(): void {
this.heartbeatTimer = setInterval(async () => {
try {
await this.redis.expire(this.options.lockKey, this.options.lockTTL);
} catch (err) {
console.error('续约 Leader 锁失败:', err);
}
}, this.options.heartbeatInterval * 1000);
}
private startMonitoring(): void {
this.monitorTimer = setInterval(async () => {
try {
const currentLeader = await this.redis.get(this.options.lockKey);
if (currentLeader === null) {
await this.tryAcquireLeadership();
}
} catch (err) {
console.error('监控 Leader 锁失败:', err);
}
}, 3000);
}
async stop(): Promise<void> {
this.isLeader = false;
if (this.heartbeatTimer) clearInterval(this.heartbeatTimer);
if (this.monitorTimer) clearInterval(this.monitorTimer);
if (this.isLeader) {
await this.releaseLock();
await this.options.onRevoked();
}
await this.redis.quit();
}
getIsLeader(): boolean {
return this.isLeader;
}
}
4.3 实例间共享状态
高可用部署中,多个索引器实例需要共享以下状态:
| 状态 | 存储位置 | 用途 |
|---|---|---|
| 当前 Leader | Redis (SETEX) | 唯一 Leader 标识 |
| 已处理高度 | PostgreSQL (indexer_state) | 断点续传 |
| 事件去重 ID | Redis (Bloom Filter) | 防止重复处理 |
| 合约最新高度 | Redis (HASH) | 分片进度追踪 |
| 心跳时间戳 | Redis (ZSET) | 健康检测 |
// src/ha/shared-state.ts
export class SharedState {
private redis: Redis;
constructor(redisUrl: string) {
this.redis = new Redis(redisUrl);
}
async setProcessedHeight(height: number): Promise<void> {
await this.redis.set('indexer:processed:height', height.toString());
}
async getProcessedHeight(): Promise<number> {
const val = await this.redis.get('indexer:processed:height');
return val ? parseInt(val) : 0;
}
async setContractProgress(contract: string, height: number): Promise<void> {
await this.redis.hset('indexer:contract:progress', contract, height.toString());
}
async getContractProgress(contract: string): Promise<number> {
const val = await this.redis.hget('indexer:contract:progress', contract);
return val ? parseInt(val) : 0;
}
private bloomKey = 'indexer:event:bloom';
async initBloom(expectedItems: number, falsePositiveRate: number): Promise<void> {
try {
await this.redis.sendCommand(
new Command('BF.RESERVE', [this.bloomKey, String(falsePositiveRate), String(expectedItems)])
);
} catch {
// BF.RESERVE 在键已存在时会报错,忽略
}
}
async isEventProcessed(eventId: string): Promise<boolean> {
const result = await this.redis.sendCommand(
new Command('BF.EXISTS', [this.bloomKey, eventId])
);
return result === 1;
}
async markEventProcessed(eventId: string): Promise<void> {
await this.redis.sendCommand(
new Command('BF.ADD', [this.bloomKey, eventId])
);
}
async reportHeartbeat(instanceId: string, ttl: number): Promise<void> {
await this.redis.setex(
`indexer:heartbeat:${instanceId}`,
ttl,
new Date().toISOString(),
);
}
async getActiveInstances(): Promise<string[]> {
const keys = await this.redis.keys('indexer:heartbeat:*');
return keys.map(k => k.replace('indexer:heartbeat:', ''));
}
}
4.4 优雅关闭与接管
[优雅关闭流程]
Leader 收到 SIGTERM
│
▼
[标记自身为 Leaving: SET indexer:leader:state "leaving"]
│
▼
[通知 Standby: PUBLISH indexer:leader:events "handover"]
│
▼
[完成当前区块处理(最多等待 30s)]
│
▼
[将 Checkpoint 写入 PostgreSQL]
│
▼
[释放 Redis 锁]
│
▼
[退出]
[Standby 接管]
Standby 收到 "handover" 通知 或 检测到锁释放
│
▼
[尝试获取 Leader 锁]
│
▼
[从最后一个 Checkpoint 恢复]
│
▼
[检查并补录 gap 区块]
│
▼
[继续正常处理]
5. 数据处理管道
5.1 管道架构
Block Source (RPC/WS)
│
▼
┌─────────────────────┐
│ Block Fetcher │ 拉取新区块 + 交易 + 事件
└────────┬────────────┘
│ Block Bundle
▼
┌─────────────────────┐
│ Event Normalizer │ 将原始 Tendermint 事件解析为标准格式
└────────┬────────────┘
│ NormalizedEvent[]
▼
┌─────────────────────┐
│ Event Filter │ 按合约地址 + 事件类型过滤
└────────┬────────────┘
│ FilteredEvent[]
▼
┌─────────────────────┐
│ State Derivation │ 从事件推导业务实体变更
│ ┌─────────────────┐│
│ │ Agent Handlers ││
│ │ PaymentHandler ││
│ │ DID Handler ││
│ │ Constitution ││
│ │ Micropayment ││
│ └─────────────────┘│
└────────┬────────────┘
│ EntityDiff[]
▼
┌─────────────────────┐
│ Reorg Detector │ 检查区块哈希一致性
└────────┬────────────┘
│ (通过则继续,不通过则回滚)
▼
┌─────────────────────┐
│ Batch Writer │ 批量写入 PostgreSQL
│ (Transactional) │
└────────┬────────────┘
│
▼
┌─────────────────────┐
│ Checkpoint Manager │ 更新已处理进度
└─────────────────────┘
5.2 Block 同步
5.2.1 WebSocket 实时同步
// src/pipeline/block-source.ts
import WebSocket from 'ws';
import { EventEmitter } from 'events';
interface BlockBundle {
height: number;
hash: string;
parentHash: string;
time: string;
proposer: string;
transactions: TxBundle[];
}
interface TxBundle {
hash: string;
height: number;
index: number;
code: number;
gasUsed: number;
gasWanted: number;
fee: string;
memo: string;
events: RawEvent[];
}
interface RawEvent {
type: string;
attributes: { key: string; value: string; index: boolean }[];
}
export class BlockSource extends EventEmitter {
private ws: WebSocket | null = null;
private restEndpoint: string;
private wsEndpoint: string;
private pollingInterval: number;
private lastHeight: number = 0;
private batchSize: number;
private pendingBlocks: BlockBundle[] = [];
constructor(config: {
wsEndpoint: string;
restEndpoint: string;
pollingInterval?: number;
batchSize?: number;
}) {
super();
this.wsEndpoint = config.wsEndpoint;
this.restEndpoint = config.restEndpoint;
this.pollingInterval = config.pollingInterval || 1000;
this.batchSize = config.batchSize || 50;
}
async startWebSocket(): Promise<void> {
this.ws = new WebSocket(this.wsEndpoint);
this.ws.on('open', () => {
this.ws!.send(JSON.stringify({
jsonrpc: '2.0',
id: 1,
method: 'subscribe',
params: ["tm.event='NewBlock'"],
}));
});
this.ws.on('message', async (data: WebSocket.Data) => {
const msg = JSON.parse(data.toString());
if (msg.method === 'Event' && msg.params?.data?.type === 'tendermint/event/NewBlock') {
const block = msg.params.data.value.block;
const height = parseInt(block.header.height);
if (height <= this.lastHeight) return;
const bundle = await this.fetchBlockBundle(height);
this.pendingBlocks.push(bundle);
if (this.pendingBlocks.length >= this.batchSize) {
this.flush();
}
}
});
this.ws.on('close', () => {
this.emit('ws_closed');
});
this.ws.on('error', (err) => {
this.emit('ws_error', err);
});
}
async startPolling(): Promise<void> {
const poll = async () => {
try {
const latestHeight = await this.fetchLatestHeight();
while (this.lastHeight < latestHeight) {
const nextHeight = this.lastHeight + 1;
const bundle = await this.fetchBlockBundle(nextHeight);
this.lastHeight = nextHeight;
this.emit('block', bundle);
}
} catch (err) {
this.emit('poll_error', err);
}
};
await poll();
setInterval(poll, this.pollingInterval);
}
private async fetchBlockBundle(height: number): Promise<BlockBundle> {
const blockResp = await fetch(
`${this.restEndpoint}/cosmos/base/tendermint/v1beta1/blocks/${height}`
);
const blockData = await blockResp.json();
const block = blockData.block;
const blockHash = blockData.block_id?.hash || '';
const txsResp = await fetch(
`${this.restEndpoint}/cosmos/tx/v1beta1/txs?events=tx.height=${height}`
);
const txsData = await txsResp.json();
const txResponses = txsData.tx_responses || [];
const transactions: TxBundle[] = txResponses.map((tx: any, index: number) => ({
hash: tx.txhash,
height,
index,
code: tx.code || 0,
gasUsed: parseInt(tx.gas_used || '0'),
gasWanted: parseInt(tx.gas_wanted || '0'),
fee: this.extractFee(tx),
memo: tx.tx?.body?.memo || '',
events: tx.events || [],
}));
return {
height,
hash: blockHash,
parentHash: block.header.last_block_id?.hash || '',
time: block.header.time,
proposer: block.header.proposer_address || '',
transactions,
};
}
private async fetchLatestHeight(): Promise<number> {
const resp = await fetch(
`${this.restEndpoint}/cosmos/base/tendermint/v1beta1/blocks/latest`
);
const data = await resp.json();
return parseInt(data.block.header.height);
}
private flush(): void {
for (const block of this.pendingBlocks) {
this.emit('block', block);
}
this.pendingBlocks = [];
}
private extractFee(tx: any): string {
const fee = tx.tx?.auth_info?.fee;
if (!fee || !fee.amount) return '0umsg';
return fee.amount.map((c: any) => `${c.amount}${c.denom}`).join(',');
}
setStartHeight(height: number): void {
this.lastHeight = height;
}
}
5.2.2 Catchup 机制
当索引器启动或落后链头时,需要通过快速追赶(Catchup)机制从断点批量拉取区块:
// src/pipeline/catchup.ts
export class CatchupManager {
private concurrency: number;
private batchSize: number;
constructor(config: { concurrency?: number; batchSize?: number }) {
this.concurrency = config.concurrency || 10;
this.batchSize = config.batchSize || 20;
}
async catchup(
fromHeight: number,
toHeight: number,
fetchBlock: (height: number) => Promise<BlockBundle>,
processBlock: (block: BlockBundle) => Promise<void>,
): Promise<void> {
const totalBlocks = toHeight - fromHeight + 1;
console.log(`Catchup 从 ${fromHeight} 到 ${toHeight},共 ${totalBlocks} 个区块`);
for (let batchStart = fromHeight; batchStart <= toHeight; batchStart += this.batchSize) {
const batchEnd = Math.min(batchStart + this.batchSize - 1, toHeight);
const promises: Promise<BlockBundle>[] = [];
for (let h = batchStart; h <= batchEnd; h++) {
promises.push(fetchBlock(h));
}
const blocks = await Promise.all(promises);
blocks.sort((a, b) => a.height - b.height);
for (const block of blocks) {
await processBlock(block);
}
console.log(`Catchup 进度: ${batchEnd}/${toHeight} (${((batchEnd / toHeight) * 100).toFixed(1)}%)`);
}
}
}
5.3 事件过滤
// src/pipeline/event-filter.ts
export interface FilterRule {
contractAddress?: string;
actions: string[];
enabled: boolean;
}
export class EventFilter {
private rules: Map<string, FilterRule> = new Map();
constructor() {
this.registerRule('agent_registry', {
contractAddress: 'msg14hj2tavq8fpeswxxa0w5xlf3v6n3a2m9v0p5k6',
actions: ['register_agent', 'update_agent', 'deregister_agent', 'set_agent_status'],
enabled: true,
});
this.registerRule('agent_payment', {
contractAddress: 'msg14hj2tavq8fpeswxxa0w5xlf3v6n3a2m9v0p5k7',
actions: ['create_session', 'fund_session', 'release_payment', 'dispute_payment', 'close_session', 'complete_milestone'],
enabled: true,
});
this.registerRule('did_registry', {
contractAddress: 'msg14hj2tavq8fpeswxxa0w5xlf3v6n3a2m9v0p5k8',
actions: ['create_did', 'update_did', 'deactivate_did', 'add_verification_method', 'remove_verification_method'],
enabled: true,
});
this.registerRule('constitution', {
contractAddress: 'msg14hj2tavq8fpeswxxa0w5xlf3v6n3a2m9v0p5k9',
actions: ['add_rule', 'remove_rule', 'update_rule', 'set_active'],
enabled: true,
});
this.registerRule('micropayment', {
contractAddress: 'msg14hj2tavq8fpeswxxa0w5xlf3v6n3a2m9v0p6k0',
actions: ['open_channel', 'deposit_channel', 'claim_channel', 'close_channel', 'update_channel_state'],
enabled: true,
});
}
registerRule(name: string, rule: FilterRule): void {
this.rules.set(name, rule);
}
filter(events: NormalizedEvent[]): NormalizedEvent[] {
return events.filter(event => {
for (const [, rule] of this.rules) {
if (!rule.enabled) continue;
const matchesContract = !rule.contractAddress ||
event.contractAddress === rule.contractAddress;
const matchesAction = rule.actions.includes(event.action);
if (matchesContract && matchesAction) return true;
}
return false;
});
}
shouldProcess(contractAddress: string, action: string): boolean {
for (const [, rule] of this.rules) {
if (!rule.enabled) continue;
const matchesContract = !rule.contractAddress || contractAddress === rule.contractAddress;
const matchesAction = rule.actions.includes(action);
if (matchesContract && matchesAction) return true;
}
return false;
}
}
5.4 状态推导
状态推导将原始事件转换为业务实体变更。以 Agent 注册事件为例:
原始事件:
type: "wasm"
attributes: {
_contract_address: "msg14hj2tavq...k6",
action: "register_agent",
agent_id: "agent-001",
owner: "msg1abc...",
name: "Agent Alpha",
agent_type: "ai",
description: "AI assistant",
capabilities: "chat,analysis",
}
状态推导:
1. 检测到 action = "register_agent"
2. 查找 agents 表: agent_id = "agent-001" 是否已存在?
3. 不存在 → INSERT (agent-001, owner, name, ...)
4. 存在 → UPDATE (更新字段)
5. 写入 agent_history 表记录变更
输出:
EntityDiff {
type: "agent",
operation: "upsert",
entityId: "agent-001",
data: { owner, name, status, ... },
}
// src/pipeline/state-derivation.ts
export interface EntityDiff {
type: 'agent' | 'session' | 'did_document' | 'constitution' | 'rule' | 'channel' | 'transaction' | 'block';
operation: 'insert' | 'update' | 'delete';
entityId: string;
data: Record<string, unknown>;
}
export class StateDerivationEngine {
private handlers: Map<string, EventHandler> = new Map();
constructor() {
this.registerHandler('register_agent', this.handleRegisterAgent);
this.registerHandler('update_agent', this.handleUpdateAgent);
this.registerHandler('deregister_agent', this.handleDeregisterAgent);
this.registerHandler('set_agent_status', this.handleSetAgentStatus);
this.registerHandler('create_session', this.handleCreateSession);
this.registerHandler('fund_session', this.handleFundSession);
this.registerHandler('release_payment', this.handleReleasePayment);
this.registerHandler('dispute_payment', this.handleDisputePayment);
this.registerHandler('close_session', this.handleCloseSession);
this.registerHandler('complete_milestone', this.handleCompleteMilestone);
this.registerHandler('create_did', this.handleCreateDID);
this.registerHandler('update_did', this.handleUpdateDID);
this.registerHandler('deactivate_did', this.handleDeactivateDID);
this.registerHandler('add_rule', this.handleAddRule);
this.registerHandler('remove_rule', this.handleRemoveRule);
this.registerHandler('update_rule', this.handleUpdateRule);
this.registerHandler('set_active', this.handleSetActive);
this.registerHandler('open_channel', this.handleOpenChannel);
this.registerHandler('deposit_channel', this.handleDeposit);
this.registerHandler('claim_channel', this.handleClaim);
this.registerHandler('close_channel', this.handleCloseChannel);
}
private registerHandler(action: string, handler: EventHandler): void {
this.handlers.set(action, handler);
}
derive(event: NormalizedEvent): EntityDiff | null {
const handler = this.handlers.get(event.action);
if (!handler) return null;
return handler(event);
}
private handleRegisterAgent(event: NormalizedEvent): EntityDiff {
return {
type: 'agent',
operation: 'insert',
entityId: event.attributes['agent_id'],
data: {
agent_id: event.attributes['agent_id'],
owner: event.attributes['owner'],
name: event.attributes['name'],
agent_type: event.attributes['agent_type'],
description: event.attributes['description'] || '',
capabilities: this.parseJsonArray(event.attributes['capabilities'] || '[]'),
status: 'ACTIVE',
created_height: event.blockHeight,
created_at: event.timestamp,
updated_at: event.timestamp,
},
};
}
private handleUpdateAgent(event: NormalizedEvent): EntityDiff {
return {
type: 'agent',
operation: 'update',
entityId: event.attributes['agent_id'],
data: this.buildUpdateData(event.attributes),
};
}
private handleDeregisterAgent(event: NormalizedEvent): EntityDiff {
return {
type: 'agent',
operation: 'update',
entityId: event.attributes['agent_id'],
data: {
status: 'DEREGISTERED',
status_reason: event.attributes['reason'] || '注销',
updated_at: event.timestamp,
},
};
}
private handleSetAgentStatus(event: NormalizedEvent): EntityDiff {
return {
type: 'agent',
operation: 'update',
entityId: event.attributes['agent_id'],
data: {
status: event.attributes['status'],
status_reason: event.attributes['reason'] || '',
updated_at: event.timestamp,
},
};
}
private handleCreateSession(event: NormalizedEvent): EntityDiff {
return {
type: 'session',
operation: 'insert',
entityId: event.attributes['session_id'],
data: {
session_id: event.attributes['session_id'],
agent_id: event.attributes['agent_id'],
client: event.attributes['client'],
amount: event.attributes['amount'],
balance: '0',
status: 'OPEN',
created_height: event.blockHeight,
created_at: event.timestamp,
updated_at: event.timestamp,
},
};
}
private handleFundSession(event: NormalizedEvent): EntityDiff {
return {
type: 'session',
operation: 'update',
entityId: event.attributes['session_id'],
data: {
balance: event.attributes['total_balance'] || event.attributes['amount'],
status: 'FUNDED',
updated_at: event.timestamp,
},
};
}
private handleReleasePayment(event: NormalizedEvent): EntityDiff {
return {
type: 'session',
operation: 'update',
entityId: event.attributes['session_id'],
data: {
status: 'RELEASED',
updated_at: event.timestamp,
},
};
}
private handleDisputePayment(event: NormalizedEvent): EntityDiff {
return {
type: 'session',
operation: 'update',
entityId: event.attributes['session_id'],
data: {
status: 'DISPUTED',
updated_at: event.timestamp,
},
};
}
private handleCloseSession(event: NormalizedEvent): EntityDiff {
return {
type: 'session',
operation: 'update',
entityId: event.attributes['session_id'],
data: {
status: event.attributes['refund_amount'] ? 'CLOSED' : 'SETTLED',
closed_at: event.timestamp,
closed_reason: event.attributes['reason'] || '',
updated_at: event.timestamp,
},
};
}
private handleCompleteMilestone(event: NormalizedEvent): EntityDiff {
return {
type: 'session',
operation: 'update',
entityId: event.attributes['session_id'],
data: {
status: 'IN_PROGRESS',
updated_at: event.timestamp,
},
};
}
private handleCreateDID(event: NormalizedEvent): EntityDiff {
return {
type: 'did_document',
operation: 'insert',
entityId: event.attributes['did'],
data: {
did: event.attributes['did'],
controller: event.attributes['controller'] || event.attributes['sender'],
public_key: event.attributes['public_key'] || '',
authentication_methods: this.parseJsonArray(event.attributes['authentication_methods'] || '[]'),
service_endpoints: this.parseJsonArray(event.attributes['service_endpoints'] || '[]'),
active: true,
created_height: event.blockHeight,
created_at: event.timestamp,
updated_at: event.timestamp,
},
};
}
private handleUpdateDID(event: NormalizedEvent): EntityDiff {
return {
type: 'did_document',
operation: 'update',
entityId: event.attributes['did'],
data: {
service_endpoints: this.parseJsonArray(event.attributes['updated_services'] || '[]'),
updated_at: event.timestamp,
},
};
}
private handleDeactivateDID(event: NormalizedEvent): EntityDiff {
return {
type: 'did_document',
operation: 'update',
entityId: event.attributes['did'],
data: {
active: false,
deactivated_at: event.timestamp,
updated_at: event.timestamp,
},
};
}
private handleAddRule(event: NormalizedEvent): EntityDiff {
return {
type: 'rule',
operation: 'insert',
entityId: event.attributes['rule_id'],
data: {
rule_id: event.attributes['rule_id'],
rule_type: event.attributes['rule_type'],
description: event.attributes['description'] || '',
parameters: event.attributes['parameters'] || '{}',
active: true,
created_height: event.blockHeight,
created_at: event.timestamp,
updated_at: event.timestamp,
},
};
}
private handleRemoveRule(event: NormalizedEvent): EntityDiff {
return {
type: 'rule',
operation: 'delete',
entityId: event.attributes['rule_id'],
data: {
removed_at: event.timestamp,
},
};
}
private handleUpdateRule(event: NormalizedEvent): EntityDiff {
return {
type: 'rule',
operation: 'update',
entityId: event.attributes['rule_id'],
data: {
updated_at: event.timestamp,
},
};
}
private handleSetActive(event: NormalizedEvent): EntityDiff {
return {
type: 'constitution',
operation: 'update',
entityId: event.attributes['constitution_id'] || 'default',
data: {
active: event.attributes['active'] === 'true',
version: event.attributes['version'] || '1.0',
updated_at: event.timestamp,
},
};
}
private handleOpenChannel(event: NormalizedEvent): EntityDiff {
return {
type: 'channel',
operation: 'insert',
entityId: event.attributes['channel_id'],
data: {
channel_id: event.attributes['channel_id'],
participant_a: event.attributes['participant_a'],
participant_b: event.attributes['participant_b'],
capacity: event.attributes['capacity'],
balance_a: '0',
balance_b: '0',
status: 'OPEN',
expiry: event.attributes['expiry'] || null,
created_height: event.blockHeight,
created_at: event.timestamp,
updated_at: event.timestamp,
},
};
}
private handleDeposit(event: NormalizedEvent): EntityDiff {
return {
type: 'channel',
operation: 'update',
entityId: event.attributes['channel_id'],
data: {
balance_a: event.attributes['total_balance_a'],
balance_b: event.attributes['total_balance_b'],
updated_at: event.timestamp,
},
};
}
private handleClaim(event: NormalizedEvent): EntityDiff {
return {
type: 'channel',
operation: 'update',
entityId: event.attributes['channel_id'],
data: {
status: 'SETTLING',
updated_at: event.timestamp,
},
};
}
private handleCloseChannel(event: NormalizedEvent): EntityDiff {
return {
type: 'channel',
operation: 'update',
entityId: event.attributes['channel_id'],
data: {
status: 'CLOSED',
final_balance_a: event.attributes['final_balance_a'],
final_balance_b: event.attributes['final_balance_b'],
closed_at: event.timestamp,
updated_at: event.timestamp,
},
};
}
private parseJsonArray(value: string): string[] {
try { return JSON.parse(value); }
catch { return value ? value.split(',').map(s => s.trim()) : []; }
}
private buildUpdateData(attrs: Record<string, string>): Record<string, unknown> {
const data: Record<string, unknown> = {};
const updatableFields = ['name', 'description', 'metadata', 'endpoint', 'fee'];
for (const field of updatableFields) {
if (attrs[field] !== undefined && attrs[field] !== '') {
data[field] = attrs[field];
}
}
if (attrs['capabilities']) {
data['capabilities'] = this.parseJsonArray(attrs['capabilities']);
}
data['updated_at'] = new Date().toISOString();
return data;
}
}
type EventHandler = (event: NormalizedEvent) => EntityDiff | null;
5.5 Reorg 处理
Cosmos SDK 基于 Tendermint BFT,理论上不会发生链重组(最终确定性)。但在实践中,以下场景仍可能导致已确认的区块发生变化:
- 节点重启后的状态同步:重启节点可能从不同的对等节点同步到不同的区块
- 分叉检测:Tendermint 的
last_commit_hash可用于检测 fork - Light Client 动态验证:轻客户端验证可能发现不一致
5.5.1 Reorg 检测
[每处理一个区块]
│
▼
[检查当前区块的 parentHash 是否等于已存储的上一个区块的 hash]
│
├── 匹配 ──→ [正常处理,推进高度]
│
└── 不匹配 ──→ [可能的 Reorg 事件]
│
▼
[向后追溯 N 个区块]
│
▼
[找到共同祖先高度 A]
│
▼
[删除 A+1 到当前高度的所有数据]
│
▼
[从 A+1 开始重新拉取和索引]
// src/pipeline/reorg-handler.ts
export class ReorgHandler {
private safeDepth: number;
constructor(safeDepth: number = 10) {
this.safeDepth = safeDepth;
}
async detectReorg(
incomingBlock: BlockBundle,
getStoredBlockHash: (height: number) => Promise<string | null>,
): Promise<ReorgResult> {
if (incomingBlock.height <= 1) {
return { detected: false, commonAncestor: incomingBlock.height };
}
const parentHash = incomingBlock.parentHash;
const storedParentHash = await getStoredBlockHash(incomingBlock.height - 1);
if (storedParentHash === null) {
return { detected: false, commonAncestor: incomingBlock.height - 1 };
}
if (parentHash === storedParentHash) {
return { detected: false, commonAncestor: incomingBlock.height - 1 };
}
const commonAncestor = await this.findCommonAncestor(
incomingBlock.height,
getStoredBlockHash,
);
return {
detected: true,
commonAncestor,
affectedHeights: Array.from(
{ length: incomingBlock.height - commonAncestor },
(_, i) => commonAncestor + 1 + i,
),
};
}
private async findCommonAncestor(
currentHeight: number,
getStoredBlockHash: (height: number) => Promise<string | null>,
): Promise<number> {
for (let depth = 1; depth <= this.safeDepth; depth++) {
const checkHeight = currentHeight - depth;
if (checkHeight <= 0) return 0;
const storedHash = await getStoredBlockHash(checkHeight);
const incomingHash = await this.getHashFromNode(checkHeight);
if (storedHash === incomingHash) {
return checkHeight;
}
}
console.warn('Reorg 超过安全深度,需要完全重新索引');
return 0;
}
private async getHashFromNode(height: number): Promise<string> {
const resp = await fetch(
`http://localhost:1317/cosmos/base/tendermint/v1beta1/blocks/${height}`
);
const data = await resp.json();
return data.block_id?.hash || '';
}
}
interface ReorgResult {
detected: boolean;
commonAncestor: number;
affectedHeights?: number[];
}
5.5.2 Reorg 数据回滚
// src/pipeline/reorg-rollback.ts
export class RollbackManager {
private pgPool: Pool;
constructor(pgPool: Pool) {
this.pgPool = pgPool;
}
async rollbackToHeight(targetHeight: number): Promise<void> {
const client = await this.pgPool.connect();
try {
await client.query('BEGIN');
await client.query('DELETE FROM events WHERE block_height > $1', [targetHeight]);
await client.query('DELETE FROM messages WHERE block_height > $1', [targetHeight]);
await client.query('DELETE FROM transactions WHERE block_height > $1', [targetHeight]);
await client.query('DELETE FROM blocks WHERE height > $1', [targetHeight]);
await client.query(
'UPDATE indexer_state SET last_processed_height = $1, updated_at = NOW() WHERE id = 1',
[targetHeight],
);
await client.query(
`UPDATE agents SET reindex_needed = true
WHERE created_height > $1`,
[targetHeight],
);
await client.query('COMMIT');
console.log(`Reorg 回滚完成: 已回退到高度 ${targetHeight}`);
} catch (err) {
await client.query('ROLLBACK');
console.error('Reorg 回滚失败:', err);
throw err;
} finally {
client.release();
}
}
}
5.6 Checkpoint 管理
// src/pipeline/checkpoint-manager.ts
export class CheckpointManager {
private checkpointInterval: number;
private lastCheckpointHeight: number = 0;
constructor(
private pgPool: Pool,
checkpointInterval: number = 100,
) {
this.checkpointInterval = checkpointInterval;
}
async getLastCheckpoint(): Promise<number> {
const result = await this.pgPool.query(
'SELECT last_processed_height FROM indexer_state WHERE id = 1'
);
if (result.rows.length === 0) {
await this.pgPool.query(
'INSERT INTO indexer_state (id, last_processed_height) VALUES (1, 0)'
);
return 0;
}
return parseInt(result.rows[0].last_processed_height);
}
async updateCheckpoint(height: number): Promise<void> {
if (height - this.lastCheckpointHeight < this.checkpointInterval) return;
await this.pgPool.query(
`UPDATE indexer_state
SET last_processed_height = $1, updated_at = NOW()
WHERE id = 1 AND last_processed_height < $1`,
[height],
);
this.lastCheckpointHeight = height;
}
async forceCheckpoint(height: number): Promise<void> {
await this.pgPool.query(
`UPDATE indexer_state
SET last_processed_height = $1, updated_at = NOW()
WHERE id = 1`,
[height],
);
this.lastCheckpointHeight = height;
}
}
6. 存储选型
6.1 存储方案对比
| 存储 | 适用场景 | 延迟 | QPS | 持久性 | 成本 |
|---|---|---|---|---|---|
| PostgreSQL 15+ | 主数据存储、关联查询、事务 | 1-10ms | 10K+ | 高(WAL + 副本) | 中 |
| TimescaleDB | 时序数据(区块时间、Gas 统计) | 1-10ms | 100K+ | 高 | 中 |
| Redis 7+ | 缓存、队列、实时去重 | <1ms | 100K+ | 中(AOF/RDB) | 低 |
| S3/MinIO | 冷归档、历史快照、备份 | 50-200ms | 1K+ | 高(多 AZ) | 极低 |
| BadgerDB (Indexer 本地) | 本地游标、临时状态 | <1ms | 50K+ | 中 | 无 |
6.2 PostgreSQL 主存储
6.2.1 表空间规划
-- 索引器表空间(SSD)
CREATE TABLESPACE indexer_ts LOCATION '/data/postgresql/indexer';
ALTER TABLE blocks SET TABLESPACE indexer_ts;
ALTER TABLE transactions SET TABLESPACE indexer_ts;
ALTER TABLE events SET TABLESPACE indexer_ts;
ALTER TABLE agents SET TABLESPACE indexer_ts;
ALTER TABLE sessions SET TABLESPACE indexer_ts;
ALTER TABLE messages SET TABLESPACE indexer_ts;
CREATE TABLESPACE archive_ts LOCATION '/data/postgresql/archive';
CREATE TABLE events_archive (
LIKE events INCLUDING ALL
) TABLESPACE archive_ts;
CREATE TABLE events_2026_07 (
CHECK (created_at >= '2026-07-01' AND created_at < '2026-08-01')
) INHERITS (events) TABLESPACE indexer_ts;
CREATE TABLE events_2026_08 (
CHECK (created_at >= '2026-08-01' AND created_at < '2026-09-01')
) INHERITS (events) TABLESPACE indexer_ts;
6.2.2 索引优化
CREATE INDEX idx_events_lookup ON events (contract_address, action, block_height DESC)
INCLUDE (tx_hash, attributes);
CREATE INDEX idx_agents_active ON agents (agent_id)
WHERE status IN ('ACTIVE', 'INACTIVE');
CREATE INDEX idx_blocks_height_brin ON blocks USING BRIN (height)
WITH (pages_per_range = 32);
CREATE INDEX idx_events_block_brin ON events USING BRIN (block_height)
WITH (pages_per_range = 32);
6.2.3 PostgreSQL 配置调优
# postgresql.conf 关键参数(32GB 内存服务器)
shared_buffers = 8GB
effective_cache_size = 24GB
work_mem = 128MB
maintenance_work_mem = 2GB
wal_buffers = 64MB
max_wal_size = 8GB
min_wal_size = 2GB
checkpoint_completion_target = 0.9
random_page_cost = 1.1
effective_io_concurrency = 200
max_parallel_workers_per_gather = 4
max_parallel_workers = 8
max_connections = 200
wal_level = replica
max_wal_senders = 4
wal_keep_size = 4096
6.3 TimescaleDB 时序存储
CREATE TABLE block_metrics (
time TIMESTAMPTZ NOT NULL,
height BIGINT NOT NULL,
tx_count INTEGER NOT NULL,
gas_used BIGINT NOT NULL,
block_time INTERVAL,
block_size INTEGER,
proposer VARCHAR(64)
);
SELECT create_hypertable('block_metrics', 'time',
chunk_time_interval => INTERVAL '1 day');
CREATE MATERIALIZED VIEW hourly_block_stats
WITH (timescaledb.continuous) AS
SELECT
time_bucket('1 hour', time) AS bucket,
COUNT(*) AS block_count,
SUM(tx_count) AS total_tx,
AVG(tx_count) AS avg_tx_per_block,
MAX(tx_count) AS max_tx_per_block,
SUM(gas_used) AS total_gas,
AVG(gas_used) AS avg_gas_per_block
FROM block_metrics
GROUP BY bucket;
SELECT add_continuous_aggregate_policy('hourly_block_stats',
start_offset => INTERVAL '3 days',
end_offset => INTERVAL '1 hour',
schedule_interval => INTERVAL '1 hour'
);
6.4 Redis 缓存与队列
// src/storage/redis-cache.ts
export class RedisCache {
private redis: Redis;
private defaultTTL: number;
constructor(redisUrl: string, defaultTTL: number = 300) {
this.redis = new Redis(redisUrl);
this.defaultTTL = defaultTTL;
}
async cacheEntity(type: string, id: string, data: unknown): Promise<void> {
const key = `entity:${type}:${id}`;
await this.redis.setex(key, this.defaultTTL, JSON.stringify(data));
}
async getCachedEntity<T>(type: string, id: string): Promise<T | null> {
const key = `entity:${type}:${id}`;
const data = await this.redis.get(key);
return data ? JSON.parse(data) : null;
}
async setContractProgress(contract: string, height: number): Promise<void> {
await this.redis.hset('indexer:progress', contract, height.toString());
}
async getContractProgress(contract: string): Promise<number> {
const val = await this.redis.hget('indexer:progress', contract);
return val ? parseInt(val) : 0;
}
async publishEvent(event: NormalizedEvent): Promise<void> {
await this.redis.publish(
'events:realtime',
JSON.stringify(event),
);
}
async cacheQueryResult(key: string, data: unknown, ttl: number): Promise<void> {
await this.redis.setex(`query:${key}`, ttl, JSON.stringify(data));
}
async getCachedQueryResult<T>(key: string): Promise<T | null> {
const data = await this.redis.get(`query:${key}`);
return data ? JSON.parse(data) : null;
}
}
6.5 S3 冷存储
// src/storage/cold-storage.ts
export class ColdStorage {
private s3: S3Client;
private bucket: string;
constructor(config: { endpoint: string; region: string; bucket: string; accessKey: string; secretKey: string }) {
this.s3 = new S3Client({
endpoint: config.endpoint,
region: config.region,
credentials: {
accessKeyId: config.accessKey,
secretAccessKey: config.secretKey,
},
forcePathStyle: true,
});
this.bucket = config.bucket;
}
async archiveBlocks(fromHeight: number, toHeight: number): Promise<void> {
const key = `blocks/${fromHeight}-${toHeight}.jsonl`;
const blocks = await this.exportBlocks(fromHeight, toHeight);
const lines = blocks.map(b => JSON.stringify(b)).join('\n');
await this.s3.send(new PutObjectCommand({
Bucket: this.bucket,
Key: key,
Body: lines,
ContentType: 'application/jsonl',
StorageClass: 'GLACIER',
}));
console.log(`已归档区块 ${fromHeight} - ${toHeight} 到 S3: ${key}`);
}
async restoreBlocks(fromHeight: number, toHeight: number): Promise<void> {
const key = `blocks/${fromHeight}-${toHeight}.jsonl`;
const response = await this.s3.send(new GetObjectCommand({
Bucket: this.bucket,
Key: key,
}));
const body = await response.Body!.transformToString();
const blocks = body.split('\n').filter(Boolean).map(line => JSON.parse(line));
await this.importBlocks(blocks);
}
private async exportBlocks(fromHeight: number, toHeight: number): Promise<unknown[]> {
return [];
}
private async importBlocks(blocks: unknown[]): Promise<void> {
}
}
6.6 存储策略总结
| 数据类型 | 存储 | 保留策略 | 查询频率 | 重要性 |
|---|---|---|---|---|
| 区块头 | PostgreSQL | 永久 | 低 | 关键 |
| 交易 | PostgreSQL | 永久 | 高 | 关键 |
| 事件 | PostgreSQL | 永久(冷数据 6 月后归档) | 高 | 关键 |
| Agent 注册 | PostgreSQL | 永久 | 极高 | 关键 |
| 支付会话 | PostgreSQL | 永久 | 极高 | 关键 |
| DID 文档 | PostgreSQL | 永久 | 中 | 关键 |
| 宪章规则 | PostgreSQL | 永久 | 低 | 关键 |
| 通道状态 | PostgreSQL | 永久 | 低 | 关键 |
| 区块指标时序 | TimescaleDB | 原始 7 天,聚合 1 年 | 低 | 辅助 |
| 热门实体 | Redis | 5 分钟 TTL | 极高 | 缓存 |
| 查询结果 | Redis | 自定义 TTL | 极高 | 缓存 |
| 历史快照 | S3/Glacier | 永久 | 极低 | 备份 |
| 本地游标 | BadgerDB | 重启即失效 | 临时 | 辅助 |
7. 索引器故障转移
7.1 故障场景分析
| 故障类型 | 影响范围 | 恢复手段 | RTO | RPO |
|---|---|---|---|---|
| Indexer 进程崩溃 | 索引暂停 | Standby 自动接管 | < 30s | < 1 block |
| Indexer 机器宕机 | 索引暂停 | Standby 自动接管 | < 60s | < 1 block |
| RPC 节点不可用 | 数据源不可达 | 切换备用 RPC 节点 | < 10s | 0 |
| PostgreSQL 主库故障 | 写入中断 | 自动故障切换至副本 | < 60s | < 1 block |
| Redis 集群降级 | 缓存 + 锁失效 | 重建缓存 + 重新选举 | < 120s | < 1 block |
| 网络分区 | 脑裂风险 | Redis 锁 + fencing | < 30s | < 1 block |
7.2 故障转移详细流程
[Primary Indexer 宕机]
│
▼
[Standby Indexer 检测到 Leader 锁过期] ← 每 3s 检查
│
▼
[Standby 获取 Redis 锁: SET indexer:leader:{chainId} {instanceId} NX EX 15]
│
├── 成功 ──→ [晋升为 Leader]
│ │
│ ▼
│ [从 PostgreSQL indexer_state 读取最后处理高度 H]
│ │
│ ▼
│ [启动 Catchup: 从 H+1 开始同步]
│ │
│ ▼
│ [验证数据完整性]
│ │
│ ▼
│ [更新 DNS/LB 指向新 Leader (可选)]
│ │
│ ▼
│ [心跳确认: 写入 indexer:heartbeat:{instanceId}]
│
└── 失败 ──→ [等待下一轮竞争]
7.3 消费位置恢复
// src/ha/resume-manager.ts
export class ResumeManager {
constructor(
private pgPool: Pool,
private checkpointMgr: CheckpointManager,
) {}
async getResumePosition(): Promise<ResumePosition> {
const pgHeight = await this.checkpointMgr.getLastCheckpoint();
return {
blockHeight: pgHeight,
fromCheckpoint: true,
};
}
async verifyAndAdjust(startHeight: number): Promise<number> {
const stored = await this.pgPool.query(
'SELECT hash FROM blocks WHERE height = $1',
[startHeight],
);
if (stored.rows.length === 0) {
const previous = await this.pgPool.query(
'SELECT height FROM blocks ORDER BY height DESC LIMIT 1'
);
return previous.rows.length > 0
? parseInt(previous.rows[0].height) + 1
: 0;
}
const currentHash = await this.fetchBlockHashFromNode(startHeight);
if (stored.rows[0].hash !== currentHash) {
console.warn(`高度 ${startHeight} 哈希不匹配,执行 Reorg 回滚`);
return -startHeight;
}
return startHeight + 1;
}
private async fetchBlockHashFromNode(height: number): Promise<string> {
const resp = await fetch(
`http://localhost:1317/cosmos/base/tendermint/v1beta1/blocks/${height}`
);
const data = await resp.json();
return data.block_id?.hash || '';
}
}
interface ResumePosition {
blockHeight: number;
fromCheckpoint: boolean;
}
7.4 脑裂防护
在 Active-Passive 模式下,网络分区可能导致两个实例都认为自己是 Leader:
[网络分区场景]
┌─────────────────────┐ ┌─────────────────────┐
│ Primary (Zone A) │ │ Primary (Zone B) │
│ 持有 Redis 锁 │←──X──→│ 检测到锁过期 │
│ 正常写 PostgreSQL │ │ 获取新锁 │
│ │ │ 也写 PostgreSQL │
└─────────────────────┘ └─────────────────────┘
│ │
▼ ▼
┌─────────────────────┐ ┌─────────────────────┐
│ PostgreSQL │ │ PostgreSQL │
│ (Primary) │ │ (Primary) ← 冲突! │
└─────────────────────┘ └─────────────────────┘
防护措施:
- Fencing Token:每次获取 Leader 锁时生成递增 token,写入数据时携带此 token,PostgreSQL 拒绝低 token 的写入
ALTER TABLE indexer_state ADD COLUMN fencing_token BIGINT NOT NULL DEFAULT 0;
UPDATE indexer_state
SET last_processed_height = $1, fencing_token = $2, updated_at = NOW()
WHERE id = 1 AND fencing_token < $2;
- STONITH(Shoot the Other Node in the Head)
// src/ha/fencing.ts
export class Fencing {
async executeFencing(instanceId: string): Promise<void> {
console.warn(`执行 Fencing: 关闭实例 ${instanceId}`);
try {
await this.sshStop(instanceId);
} catch {
console.warn(`SSH 关闭 ${instanceId} 失败,等待锁自动过期`);
}
await this.forceReleaseLock(instanceId);
}
private async forceReleaseLock(instanceId: string): Promise<void> {
const script = `
if redis.call("GET", KEYS[1]) == ARGV[1] then
return redis.call("DEL", KEYS[1])
end
return 0
`;
await redis.eval(script, 1, 'indexer:leader:lock', instanceId);
}
private async sshStop(instanceId: string): Promise<void> {
}
}
7.5 多活数据中心
对于跨 AZ 或跨地域部署,推荐使用 Active-Passive 跨 AZ 模式:
┌─────────── Region A (Primary) ──────────────────┐
│ AZ-1a AZ-1b │
│ Indexer Primary Indexer Standby │
│ PG Primary PG Replica │
│ Redis Master Redis Replica │
└──────────────────────────────────────────────────┘
│
Async 复制 (跨 AZ)
│
┌─────────── Region B (Disaster Recovery) ─────────┐
│ AZ-2a AZ-2b │
│ Indexer Standby Indexer Standby │
│ PG Replica PG Replica │
│ Redis Replica Redis Replica │
└──────────────────────────────────────────────────┘
8. 数据一致性
8.1 一致性模型
| 模型 | 说明 | 适用场景 |
|---|---|---|
| 强一致 | 写入确认后,所有读取立即看到最新版本 | 支付会话状态、Agent 状态变更 |
| 最终一致 | 写入后延迟(通常 < 1s)才在所有副本可见 | GraphQL 查询、列表展示 |
| 读取自身写入 | 写入者能立即看到自己的写入,其他用户可能看到旧数据 | 前端即时反馈 |
8.2 MSG Chain 索引器一致性策略
Indexer → PostgreSQL Primary (强一致写入)
│
├──→ PostgreSQL Replica (异步复制,最终一致)
│ │
│ └──→ GraphQL API Read (最终一致读)
│
├──→ Redis Cache (TTL 后失效)
│ │
│ └──→ GraphQL API Cached Read (最终一致读)
│
└──→ 业务实体直接查询 Primary (强一致读)
│
└──→ Agent API (关键操作)
8.3 确认延迟分析
Block Time → RPC Propagation → Indexer Processing → DB Write → API Available
MSG Chain 出块时间: ~3s
RPC 传播延迟: 100-500ms
Indexer 处理延迟: 50-200ms
DB 写入延迟: 5-20ms
API 查询延迟: 1-10ms
端到端延迟: 3.2s - 3.7s
8.4 一致性保证策略
// src/consistency/consistency-guard.ts
export class ConsistencyGuard {
private confirmedLag: number = 3;
async waitForConfirmation(
targetHeight: number,
getCurrentHeight: () => Promise<number>,
): Promise<boolean> {
while (true) {
const currentHeight = await getCurrentHeight();
if (currentHeight >= targetHeight + this.confirmedLag) {
return true;
}
await new Promise(resolve => setTimeout(resolve, 1000));
}
}
async readFromReplica<T>(
query: string,
params: unknown[],
): Promise<T> {
return this.replicaPool.query(query, params);
}
async readFromPrimary<T>(
query: string,
params: unknown[],
): Promise<T> {
return this.primaryPool.query(query, params);
}
async readWithConsistency<T>(
query: string,
params: unknown[],
options: { strong?: boolean; criticalType?: string },
): Promise<T> {
if (options.criticalType === 'payment' || options.strong) {
return this.readFromPrimary(query, params);
}
return this.readFromReplica(query, params);
}
}
8.5 延迟监控公式
确认延迟 (s) = (当前链头高度 - 索引器最后处理高度) × 平均出块时间 + 管道处理时间
当前链头高度: 从节点获取
索引器最后处理高度: 从 indexer_state 表读取
平均出块时间: MSG Chain 约 3s
管道处理时间: 50-200ms (少量区块) / 1-5s (大量 catchup 时)
9. 性能基准
9.1 TPS 基准
基于 MSG Chain 的 CosmWasm 合约特性,索引器的性能基准如下:
| 组件 | 指标 | 测量值 | 条件 |
|---|---|---|---|
| Block Source | 区块拉取 | 50-100 blocks/s | WebSocket + REST 混合 |
| Event Parser | 事件解析 | 5,000 events/s | 单核 |
| State Derivation | 状态推导 | 3,000 events/s | 单核,含 DB 写入 |
| Batch Writer | 批量写入 | 10,000 rows/s | PostgreSQL, 1000 条一批 |
| GraphQL Query | 简单查询 | 5,000 QPS | Redis 缓存命中 |
| GraphQL Query | 复杂 JOIN | 500 QPS | PostgreSQL 直接查询 |
| Database Read | 主键查询 | 10,000 QPS | 索引命中 |
| Database Read | 范围扫描 | 2,000 QPS | 覆盖索引 |
9.2 存储增长率估算
-- 区块表
每个区块行大小: ~200 bytes
年增长率: (365 × 24 × 3600 / 3s) × 200B ≈ 2.1 GB/年
-- 交易表
每笔交易行大小: ~500 bytes
假设平均每区块 50 笔交易
年增长率: (365 × 24 × 3600 / 3s) × 50 × 500B ≈ 262 GB/年
-- 事件表
每个事件行大小: ~400 bytes (含 JSONB)
假设平均每笔交易产生 5 个事件
年增长率: (365 × 24 × 3600 / 3s) × 50 × 5 × 400B ≈ 1.05 TB/年
-- 业务实体表(agents, sessions, did_documents 等)
每个 Agent 行大小: ~1KB
假设年注册 100 万 Agent
年增长率: 1万 × 1KB ≈ 10 MB/年
每个 Session 行大小: ~800B
假设年创建 1000 万 Session
年增长率: 1000万 × 800B ≈ 8 GB/年
9.3 存储总容量估算公式
总存储 (TB/年) =
blocks_size +
transactions_size +
events_size +
entity_tables_size +
indexes_size (约为数据量的 30-50%) +
WAL_size (约为数据量的 20%) +
temp_files + connection_overhead
简化估算: events_size × 1.7 ~ 最终总存储
9.4 查询延迟基准
-- 测试 1: 按合约地址查事件
EXPLAIN ANALYZE
SELECT * FROM events
WHERE contract_address = 'msg14hj2tavq8fpeswxxa0w5xlf3v6n3a2m9v0p5k6'
AND action = 'register_agent'
ORDER BY block_height DESC
LIMIT 100;
-- 预期: 0.5-2ms
-- 测试 2: Agent 及其支付会话
EXPLAIN ANALYZE
SELECT a.agent_id, a.name, a.status,
s.session_id, s.status as session_status, s.amount
FROM agents a
LEFT JOIN sessions s ON a.agent_id = s.agent_id
WHERE a.owner = 'msg1abc123...'
ORDER BY s.created_at DESC;
-- 预期: 2-10ms
-- 测试 3: 按时间范围查交易
EXPLAIN ANALYZE
SELECT * FROM transactions
WHERE block_height BETWEEN 100000 AND 101000
ORDER BY block_height DESC;
-- 预期: 1-5ms
-- 测试 4: 全文搜索
EXPLAIN ANALYZE
SELECT * FROM transactions
WHERE search_vector @@ to_tsquery('simple', 'agent')
ORDER BY block_height DESC
LIMIT 50;
-- 预期: 10-50ms
9.5 资源需求估算
// src/performance/resource-estimator.ts
export function estimateResources(params: {
avgBlockTime: number;
avgTxPerBlock: number;
avgEventsPerTx: number;
totalContracts: number;
retentionDays: number;
peakQPS: number;
}): ResourceEstimate {
const blocksPerDay = 86400 / params.avgBlockTime;
const txsPerDay = blocksPerDay * params.avgTxPerBlock;
const eventsPerDay = txsPerDay * params.avgEventsPerTx;
const pipelineCPU = Math.ceil(eventsPerDay / (86400 * 3000) * 2);
const apiCPU = Math.ceil(params.peakQPS / 2000);
const totalCPU = Math.max(pipelineCPU + apiCPU, 4);
const pgMemory = Math.max(
Math.ceil(eventsPerDay * 400 / 1024 / 1024 * 0.2),
4,
);
const redisMemory = Math.max(
Math.ceil(params.peakQPS * params.totalContracts * 0.001),
2,
);
const totalMemory = pgMemory + redisMemory + 2;
const txSize = 500;
const eventSize = 400;
const blockSize = 200;
const dataPerDay = (
blocksPerDay * blockSize +
txsPerDay * txSize +
eventsPerDay * eventSize
) / 1024 / 1024 / 1024;
const totalDisk = Math.ceil(dataPerDay * params.retentionDays * 1.5);
return {
cpu: totalCPU,
memory: totalMemory,
disk: totalDisk,
dataPerDay,
estimatedMonthlyCost: {
cpu: totalCPU,
memory: `${totalMemory}GB`,
disk: `${totalDisk}GB`,
},
};
}
interface ResourceEstimate {
cpu: number;
memory: number;
disk: number;
dataPerDay: number;
estimatedMonthlyCost: Record<string, string>;
}
10. 监控与告警
10.1 监控指标体系
| 类别 | 指标 | 说明 | 告警阈值 | 严重程度 |
|---|---|---|---|---|
| 同步 | indexer_sync_lag_blocks | 索引器落后链头区块数 | > 10 | Warning |
| indexer_sync_lag_blocks | > 100 | Critical | ||
| indexer_sync_lag_seconds | 索引器落后时长(秒) | > 30s | Warning | |
| indexer_last_processed_height | 最后处理高度 | < 链头高度 | Info | |
| 性能 | indexer_events_per_second | 每秒处理事件数 | < 100 | Warning |
| indexer_block_processing_time | 单区块处理时间(ms) | > 1000 | Warning | |
| db_write_batch_size | 批量写入大小 | < 10 | Warning | |
| 存储 | pg_database_size | 数据库大小 | > 80% 磁盘 | Warning |
| pg_table_bloat_ratio | 表膨胀率 | > 30% | Warning | |
| index_size_growth_rate | 索引增长率 | 异常增长 | Warning | |
| 查询 | graphql_query_p99_latency | 查询 P99 延迟 | > 500ms | Warning |
| graphql_query_error_rate | 查询错误率 | > 1% | Critical | |
| api_requests_per_second | API 请求速率 | 自定义 | Info | |
| 系统 | cpu_usage | CPU 使用率 | > 80% | Warning |
| memory_usage | 内存使用率 | > 90% | Critical | |
| process_alive | 进程存活 | 0 | Critical |
10.2 监控指标暴露(Prometheus 格式)
// src/monitoring/metrics.ts
import prometheus from 'prom-client';
const registry = new prometheus.Registry();
prometheus.collectDefaultMetrics({ register: registry });
export const indexerSyncLag = new prometheus.Gauge({
name: 'indexer_sync_lag_blocks',
help: 'Indexer 落后链头区块数',
registers: [registry],
});
export const lastProcessedHeight = new prometheus.Gauge({
name: 'indexer_last_processed_height',
help: 'Indexer 最后处理的区块高度',
registers: [registry],
});
export const chainHeadHeight = new prometheus.Gauge({
name: 'chain_head_height',
help: 'MSG Chain 链头高度',
registers: [registry],
});
export const eventsProcessed = new prometheus.Counter({
name: 'indexer_events_processed_total',
help: '累计处理事件数',
labelNames: ['contract', 'action'],
registers: [registry],
});
export const blockProcessingTime = new prometheus.Histogram({
name: 'indexer_block_processing_seconds',
help: '区块处理时间分布',
buckets: [0.01, 0.05, 0.1, 0.5, 1, 2, 5],
registers: [registry],
});
export const dbWriteLatency = new prometheus.Histogram({
name: 'indexer_db_write_seconds',
help: '数据库写入延迟分布',
buckets: [0.001, 0.005, 0.01, 0.05, 0.1, 0.5],
registers: [registry],
});
export const pipelineStageDuration = new prometheus.Histogram({
name: 'indexer_pipeline_stage_seconds',
help: '管道各阶段耗时分布',
labelNames: ['stage'],
buckets: [0.001, 0.005, 0.01, 0.05, 0.1, 0.5],
registers: [registry],
});
export const reorgEvents = new prometheus.Counter({
name: 'indexer_reorg_events_total',
help: 'Reorg 事件累计数',
registers: [registry],
});
export const failoverEvents = new prometheus.Counter({
name: 'indexer_failover_events_total',
help: '故障转移事件累计数',
labelNames: ['from_instance', 'to_instance'],
registers: [registry],
});
export const leaderStatus = new prometheus.Gauge({
name: 'indexer_leader_status',
help: '当前实例是否为 Leader (1=是, 0=否)',
registers: [registry],
});
export const queryLatency = new prometheus.Histogram({
name: 'graphql_query_duration_seconds',
help: 'GraphQL 查询延迟分布(按操作名)',
labelNames: ['operation'],
buckets: [0.01, 0.05, 0.1, 0.3, 0.5, 1, 3],
registers: [registry],
});
export const queryErrors = new prometheus.Counter({
name: 'graphql_query_errors_total',
help: 'GraphQL 查询错误数',
labelNames: ['operation', 'error_code'],
registers: [registry],
});
export const databaseSize = new prometheus.Gauge({
name: 'pg_database_size_bytes',
help: 'PostgreSQL 数据库大小',
labelNames: ['database'],
registers: [registry],
});
export { registry };
10.3 告警规则配置(Prometheus AlertManager)
# alertmanager/alert-rules.yml
groups:
- name: indexer_alerts
rules:
- alert: IndexerSyncLagHigh
expr: indexer_sync_lag_blocks > 10
for: 1m
labels:
severity: warning
annotations:
summary: "索引器同步延迟较高"
description: "索引器落后链头 {{ $value }} 个区块"
- alert: IndexerSyncLagCritical
expr: indexer_sync_lag_blocks > 100
for: 1m
labels:
severity: critical
annotations:
summary: "索引器同步延迟严重"
description: "索引器落后链头 {{ $value }} 个区块,需要立即处理"
- alert: IndexerDown
expr: up{job="indexer"} == 0
for: 30s
labels:
severity: critical
annotations:
summary: "索引器进程宕机"
description: "索引器 {{ $labels.instance }} 已停止运行超过 30 秒"
- alert: NoLeaderElected
expr: indexer_leader_status == 0
for: 30s
labels:
severity: critical
annotations:
summary: "没有 Leader 实例"
description: "所有索引器实例均未获得 Leader 锁"
- alert: BlockProcessingSlow
expr: indexer_block_processing_seconds{quantile="0.95"} > 2
for: 5m
labels:
severity: warning
annotations:
summary: "区块处理过慢"
description: "P95 区块处理时间 {{ $value }}s"
- alert: ReorgDetected
expr: rate(indexer_reorg_events_total[5m]) > 0
for: 1m
labels:
severity: critical
annotations:
summary: "检测到链重组"
description: "过去 5 分钟发生了 {{ $value }} 次链重组事件"
- alert: DiskSpaceWarning
expr: (pg_database_size_bytes / node_filesystem_size_bytes{mountpoint="/data"}) * 100 > 80
for: 5m
labels:
severity: warning
annotations:
summary: "磁盘空间不足 80%"
description: "数据库已使用 {{ $value }}% 的磁盘空间"
- alert: DiskSpaceCritical
expr: (pg_database_size_bytes / node_filesystem_size_bytes{mountpoint="/data"}) * 100 > 90
for: 5m
labels:
severity: critical
annotations:
summary: "磁盘空间不足 90%"
description: "数据库已使用 {{ $value }}% 的磁盘空间,即将写满"
- alert: QueryLatencyHigh
expr: graphql_query_duration_seconds{quantile="0.99"} > 1
for: 5m
labels:
severity: warning
annotations:
summary: "查询延迟过高"
description: "P99 查询延迟 {{ $value }}s"
- alert: QueryErrorRateHigh
expr: rate(graphql_query_errors_total[5m]) / rate(graphql_queries_total[5m]) > 0.01
for: 5m
labels:
severity: critical
annotations:
summary: "查询错误率过高"
description: "错误率 {{ $value | humanizePercentage }}"
- alert: FailoverOccurred
expr: rate(indexer_failover_events_total[5m]) > 0
for: 1m
labels:
severity: info
annotations:
summary: "索引器发生故障转移"
description: "从 {{ $labels.from_instance }} 转移到 {{ $labels.to_instance }}"
10.4 落后者检测
// src/monitoring/lag-detector.ts
export class LagDetector {
private alertHistory: Map<string, number> = new Map();
constructor(
private alertThreshold: number = 10,
private criticalThreshold: number = 100,
) {}
async detectLag(lastProcessed: number, chainHead: number): Promise<LagLevel> {
const lag = chainHead - lastProcessed;
if (lag > this.criticalThreshold) return 'CRITICAL';
if (lag > this.alertThreshold) return 'WARNING';
return 'OK';
}
async reportLag(
instanceId: string,
lastProcessed: number,
chainHead: number,
): Promise<void> {
const lag = chainHead - lastProcessed;
const level = await this.detectLag(lastProcessed, chainHead);
console.log(
`[LagDetector] ${instanceId}: last=${lastProcessed} head=${chainHead} lag=${lag} level=${level}`
);
if (level === 'CRITICAL') {
const lastAlert = this.alertHistory.get(instanceId) || 0;
if (Date.now() - lastAlert > 300000) {
this.alertHistory.set(instanceId, Date.now());
this.sendAlert(instanceId, lag, level);
}
}
if (level === 'WARNING') {
const lastAlert = this.alertHistory.get(instanceId) || 0;
if (Date.now() - lastAlert > 600000) {
this.alertHistory.set(instanceId, Date.now());
this.sendAlert(instanceId, lag, level);
}
}
}
private sendAlert(instanceId: string, lag: number, level: LagLevel): void {
const message = {
title: `[${level}] 索引器 ${instanceId} 同步落后`,
message: `落后 ${lag} 个区块`,
severity: level,
timestamp: new Date().toISOString(),
instance: instanceId,
lag,
};
console.warn('告警:', JSON.stringify(message));
}
}
type LagLevel = 'OK' | 'WARNING' | 'CRITICAL';
10.5 健康检查端点
// src/monitoring/health-check.ts
import express from 'express';
export function createHealthRouter(checkers: HealthChecker[]): express.Router {
const router = express.Router();
router.get('/health', async (req, res) => {
res.json({ status: 'ok', timestamp: new Date().toISOString() });
});
router.get('/health/ready', async (req, res) => {
const results = await Promise.all(
checkers.map(c => c.check())
);
const allHealthy = results.every(r => r.healthy);
const statusCode = allHealthy ? 200 : 503;
res.status(statusCode).json({
status: allHealthy ? 'ready' : 'not_ready',
checks: results.map(r => ({
name: r.name,
healthy: r.healthy,
message: r.message,
duration: r.duration,
})),
});
});
return router;
}
interface HealthCheckResult {
name: string;
healthy: boolean;
message: string;
duration: number;
}
interface HealthChecker {
check(): Promise<HealthCheckResult>;
}
export class DatabaseHealthChecker implements HealthChecker {
constructor(private pool: Pool) {}
async check(): Promise<HealthCheckResult> {
const start = Date.now();
try {
await this.pool.query('SELECT 1');
return {
name: 'postgresql',
healthy: true,
message: '数据库连接正常',
duration: Date.now() - start,
};
} catch (err) {
return {
name: 'postgresql',
healthy: false,
message: `数据库连接失败: ${err}`,
duration: Date.now() - start,
};
}
}
}
export class RedisHealthChecker implements HealthChecker {
constructor(private redis: Redis) {}
async check(): Promise<HealthCheckResult> {
const start = Date.now();
try {
await this.redis.ping();
return {
name: 'redis',
healthy: true,
message: 'Redis 连接正常',
duration: Date.now() - start,
};
} catch (err) {
return {
name: 'redis',
healthy: false,
message: `Redis 连接失败: ${err}`,
duration: Date.now() - start,
};
}
}
}
export class ChainHealthChecker implements HealthChecker {
constructor(private rpcEndpoint: string) {}
async check(): Promise<HealthCheckResult> {
const start = Date.now();
try {
const resp = await fetch(`${this.rpcEndpoint}/status`);
const data = await resp.json();
return {
name: 'chain_rpc',
healthy: true,
message: `链 ${data.result.node_info.network} 连接正常`,
duration: Date.now() - start,
};
} catch (err) {
return {
name: 'chain_rpc',
healthy: false,
message: `RPC 节点连接失败: ${err}`,
duration: Date.now() - start,
};
}
}
}
export class SyncHealthChecker implements HealthChecker {
constructor(
private getProcessedHeight: () => Promise<number>,
private getChainHead: () => Promise<number>,
private maxLag: number,
) {}
async check(): Promise<HealthCheckResult> {
const start = Date.now();
try {
const processed = await this.getProcessedHeight();
const head = await this.getChainHead();
const lag = head - processed;
return {
name: 'sync_lag',
healthy: lag <= this.maxLag,
message: `落后 ${lag} 个区块 (阈值: ${this.maxLag})`,
duration: Date.now() - start,
};
} catch (err) {
return {
name: 'sync_lag',
healthy: false,
message: `同步检查失败: ${err}`,
duration: Date.now() - start,
};
}
}
}
11. 部署实践
11.1 Docker Compose 部署(开发/预发环境)
# docker-compose.yml
version: "3.9"
x-common: &common
restart: unless-stopped
networks:
- msgchain
services:
postgres:
image: postgres:15-alpine
<<: *common
ports:
- "5432:5432"
environment:
POSTGRES_USER: indexer
POSTGRES_PASSWORD: indexer_pass
POSTGRES_DB: indexer
POSTGRES_INITDB_ARGS: "-E UTF8 --locale=C"
volumes:
- pg_data:/var/lib/postgresql/data
- ./migrations:/docker-entrypoint-initdb.d
healthcheck:
test: ["CMD-SHELL", "pg_isready -U indexer"]
interval: 10s
timeout: 5s
retries: 5
redis:
image: redis:7-alpine
<<: *common
ports:
- "6379:6379"
volumes:
- redis_data:/data
healthcheck:
test: ["CMD", "redis-cli", "ping"]
interval: 10s
timeout: 5s
retries: 5
indexer-primary:
build:
context: .
dockerfile: Dockerfile.indexer
<<: *common
ports:
- "9464:9464"
environment:
NODE_ENV: production
INDEXER_MODE: leader
INDEXER_INSTANCE_ID: "indexer-primary"
INDEXER_CHAIN_ID: "msg-chain-1"
INDEXER_RPC_ENDPOINT: "http://msgchain-node:26657"
INDEXER_REST_ENDPOINT: "http://msgchain-node:1317"
INDEXER_WS_ENDPOINT: "ws://msgchain-node:26657/websocket"
DATABASE_URL: "postgres://indexer:indexer_pass@postgres:5432/indexer"
REDIS_URL: "redis://redis:6379"
INDEXER_START_HEIGHT: "0"
depends_on:
postgres:
condition: service_healthy
redis:
condition: service_healthy
healthcheck:
test: ["CMD", "curl", "-f", "http://localhost:9464/health"]
interval: 30s
timeout: 10s
retries: 5
indexer-standby:
build:
context: .
dockerfile: Dockerfile.indexer
<<: *common
ports:
- "9465:9464"
environment:
NODE_ENV: production
INDEXER_MODE: standby
INDEXER_INSTANCE_ID: "indexer-standby"
INDEXER_CHAIN_ID: "msg-chain-1"
INDEXER_RPC_ENDPOINT: "http://msgchain-node:26657"
INDEXER_REST_ENDPOINT: "http://msgchain-node:1317"
INDEXER_WS_ENDPOINT: "ws://msgchain-node:26657/websocket"
DATABASE_URL: "postgres://indexer:indexer_pass@postgres:5432/indexer"
REDIS_URL: "redis://redis:6379"
INDEXER_START_HEIGHT: "0"
depends_on:
postgres:
condition: service_healthy
redis:
condition: service_healthy
profiles:
- production
graphql-api:
build:
context: .
dockerfile: Dockerfile.api
<<: *common
ports:
- "4000:4000"
environment:
NODE_ENV: production
DATABASE_URL: "postgres://indexer:indexer_pass@postgres:5432/indexer"
REDIS_URL: "redis://redis:6379"
GRAPHQL_PLAYGROUND: "false"
depends_on:
postgres:
condition: service_healthy
healthcheck:
test: ["CMD", "curl", "-f", "http://localhost:4000/health"]
interval: 30s
timeout: 10s
retries: 5
prometheus:
image: prom/prometheus:v2.50.0
<<: *common
ports:
- "9090:9090"
volumes:
- ./prometheus:/etc/prometheus
- prometheus_data:/prometheus
command:
- "--config.file=/etc/prometheus/prometheus.yml"
- "--storage.tsdb.retention.time=30d"
grafana:
image: grafana/grafana:10.3.0
<<: *common
ports:
- "3000:3000"
environment:
GF_SECURITY_ADMIN_PASSWORD: admin
GF_INSTALL_PLUGINS: "grafana-piechart-panel"
volumes:
- grafana_data:/var/lib/grafana
- ./grafana/dashboards:/etc/grafana/provisioning/dashboards
- ./grafana/datasources:/etc/grafana/provisioning/datasources
volumes:
pg_data:
redis_data:
prometheus_data:
grafana_data:
networks:
msgchain:
driver: bridge
11.2 K8s 部署蓝图
# k8s/indexer-deployment.yaml
apiVersion: apps/v1
kind: Deployment
metadata:
name: indexer
namespace: msgchain
labels:
app: indexer
component: indexer
spec:
replicas: 1
selector:
matchLabels:
app: indexer
template:
metadata:
labels:
app: indexer
annotations:
prometheus.io/scrape: "true"
prometheus.io/port: "9464"
spec:
containers:
- name: indexer
image: msgchain/indexer:1.0.0
imagePullPolicy: Always
ports:
- containerPort: 9464
name: metrics
env:
- name: INDEXER_MODE
valueFrom:
fieldRef:
fieldPath: metadata.labels['indexer-mode']
- name: INDEXER_INSTANCE_ID
valueFrom:
fieldRef:
fieldPath: metadata.name
- name: INDEXER_CHAIN_ID
value: "msg-chain-1"
- name: INDEXER_RPC_ENDPOINT
value: "http://msgchain-rpc:26657"
- name: INDEXER_REST_ENDPOINT
value: "http://msgchain-rest:1317"
- name: INDEXER_WS_ENDPOINT
value: "ws://msgchain-rpc:26657/websocket"
- name: DATABASE_URL
valueFrom:
secretKeyRef:
name: indexer-secrets
key: database-url
- name: REDIS_URL
valueFrom:
secretKeyRef:
name: indexer-secrets
key: redis-url
- name: INDEXER_START_HEIGHT
value: "0"
- name: INDEXER_LOG_LEVEL
value: "info"
resources:
requests:
cpu: "2"
memory: "4Gi"
limits:
cpu: "4"
memory: "8Gi"
livenessProbe:
httpGet:
path: /health
port: 9464
initialDelaySeconds: 30
periodSeconds: 30
timeoutSeconds: 10
failureThreshold: 3
readinessProbe:
httpGet:
path: /health/ready
port: 9464
initialDelaySeconds: 15
periodSeconds: 15
timeoutSeconds: 5
failureThreshold: 2
volumeMounts:
- name: tmp
mountPath: /tmp
volumes:
- name: tmp
emptyDir: {}
---
# k8s/indexer-standby-deployment.yaml
apiVersion: apps/v1
kind: Deployment
metadata:
name: indexer-standby
namespace: msgchain
labels:
app: indexer
component: indexer
indexer-mode: standby
spec:
replicas: 1
selector:
matchLabels:
app: indexer-standby
template:
metadata:
labels:
app: indexer-standby
indexer-mode: standby
spec:
containers:
- name: indexer
image: msgchain/indexer:1.0.0
ports:
- containerPort: 9464
name: metrics
env:
- name: INDEXER_MODE
value: "standby"
---
# k8s/graphql-api-deployment.yaml
apiVersion: apps/v1
kind: Deployment
metadata:
name: graphql-api
namespace: msgchain
labels:
app: graphql-api
spec:
replicas: 2
selector:
matchLabels:
app: graphql-api
template:
metadata:
labels:
app: graphql-api
spec:
containers:
- name: graphql
image: msgchain/graphql-api:1.0.0
ports:
- containerPort: 4000
name: graphql
env:
- name: DATABASE_URL
valueFrom:
secretKeyRef:
name: indexer-secrets
key: database-url-readonly
- name: REDIS_URL
valueFrom:
secretKeyRef:
name: indexer-secrets
key: redis-url
resources:
requests:
cpu: "500m"
memory: "1Gi"
limits:
cpu: "2"
memory: "4Gi"
livenessProbe:
httpGet:
path: /health
port: 4000
readinessProbe:
httpGet:
path: /health/ready
port: 4000
---
# k8s/service.yaml
apiVersion: v1
kind: Service
metadata:
name: indexer-metrics
namespace: msgchain
labels:
app: indexer
spec:
selector:
app: indexer
ports:
- name: metrics
port: 9464
targetPort: 9464
---
apiVersion: v1
kind: Service
metadata:
name: graphql-api
namespace: msgchain
labels:
app: graphql-api
spec:
selector:
app: graphql-api
ports:
- name: graphql
port: 80
targetPort: 4000
type: ClusterIP
---
# k8s/hpa.yaml
apiVersion: autoscaling/v2
kind: HorizontalPodAutoscaler
metadata:
name: graphql-api-hpa
namespace: msgchain
spec:
scaleTargetRef:
apiVersion: apps/v1
kind: Deployment
name: graphql-api
minReplicas: 2
maxReplicas: 10
metrics:
- type: Resource
resource:
name: cpu
target:
type: Utilization
averageUtilization: 70
- type: Resource
resource:
name: memory
target:
type: Utilization
averageUtilization: 80
11.3 Helm Chart 示例
# helm/msgchain-indexer/values.yaml
global:
chainId: "msg-chain-1"
domain: msgchain.org
bech32Prefix: msg
indexer:
image:
repository: msgchain/indexer
tag: 1.0.0
pullPolicy: Always
replicas: 1
mode: "leader"
resources:
requests:
cpu: "2"
memory: "4Gi"
limits:
cpu: "4"
memory: "8Gi"
config:
rpcEndpoint: "http://msgchain-rpc:26657"
restEndpoint: "http://msgchain-rest:1317"
wsEndpoint: "ws://msgchain-rpc:26657/websocket"
startHeight: 0
logLevel: "info"
batchSize: 100
checkpointInterval: 100
reorgSafeDepth: 10
syncMode: "hybrid"
ha:
enabled: true
leaderElection:
lockKey: "indexer:leader:lock"
lockTTL: 15
heartbeatInterval: 5
fencing:
enabled: true
monitoring:
metricsPort: 9464
healthEndpoint: "/health"
readinessEndpoint: "/health/ready"
standby:
enabled: true
replicas: 1
resources:
requests:
cpu: "2"
memory: "4Gi"
limits:
cpu: "4"
memory: "8Gi"
graphql:
image:
repository: msgchain/graphql-api
tag: 1.0.0
pullPolicy: Always
replicas: 2
resources:
requests:
cpu: "500m"
memory: "1Gi"
limits:
cpu: "2"
memory: "4Gi"
config:
port: 4000
playground: false
cacheTTL: 300
postgresql:
enabled: true
primary:
resources:
requests:
cpu: "2"
memory: "8Gi"
limits:
cpu: "8"
memory: "32Gi"
persistence:
size: 500Gi
storageClass: ssd
readReplicas:
enabled: true
count: 1
resources:
requests:
cpu: "2"
memory: "8Gi"
redis:
enabled: true
cluster:
enabled: true
masterCount: 3
replicaCount: 1
resources:
requests:
cpu: "1"
memory: "2Gi"
monitoring:
prometheus:
enabled: true
retentionDays: 30
grafana:
enabled: true
adminPassword: "admin"
alertmanager:
enabled: true
slackWebhook: ""
ingress:
enabled: true
host: "indexer.msgchain.org"
tls:
enabled: true
secretName: "msgchain-tls"
11.4 网络与安全配置
# k8s/network-policy.yaml
apiVersion: networking.k8s.io/v1
kind: NetworkPolicy
metadata:
name: indexer-network-policy
namespace: msgchain
spec:
podSelector:
matchLabels:
app: indexer
policyTypes:
- Ingress
- Egress
ingress:
- from:
- podSelector:
matchLabels:
app: graphql-api
- podSelector:
matchLabels:
app: prometheus
ports:
- port: 9464
protocol: TCP
egress:
- to:
- podSelector:
matchLabels:
app: postgresql
ports:
- port: 5432
- to:
- podSelector:
matchLabels:
app: redis
ports:
- port: 6379
- to:
- ipBlock:
cidr: 10.0.0.0/8
ports:
- port: 26657
protocol: TCP
- port: 1317
protocol: TCP
12. 与 MSG Chain Indexer 数据平面的集成
12.1 数据平面架构
MSG Chain 的 Indexer 数据平面包含多个子系统,链下高可用索引器是其中的核心组件:
┌─────────────────────────────────────────────────────────────────────┐
│ MSG Chain Indexer 数据平面 │
│ │
│ ┌──────────────┐ ┌──────────────┐ ┌──────────────┐ │
│ │ The Graph │ │ 事件索引器 │ │ 链下索引器 │ │
│ │ 子图索引器 │ │ (GraphQL) │ │ (HA 部署) │ ← 本文 │
│ └──────┬───────┘ └──────┬───────┘ └──────┬───────┘ │
│ │ │ │ │
│ └─────────────────┼─────────────────┘ │
│ │ │
│ ▼ │
│ ┌────────────────────────┐ │
│ │ Indexer 数据总线 │ │
│ │ (Redis Pub/Sub + │ │
│ │ PostgreSQL Logical │ │
│ │ Replication) │ │
│ └────────┬───────────────┘ │
│ │ │
│ ┌─────────────┼─────────────┐ │
│ │ │ │ │
│ ▼ ▼ ▼ │
│ ┌───────────┐ ┌───────────┐ ┌───────────┐ │
│ │ Agent API │ │ 区块浏览器 │ │ 分析平台 │ │
│ │ 服务 │ │ 前端 │ │ (Grafana) │ │
│ └───────────┘ └───────────┘ └───────────┘ │
└─────────────────────────────────────────────────────────────────────┘
12.2 集成点
| 集成目标 | 方式 | 数据流动 | 延迟要求 |
|---|---|---|---|
| The Graph 子图 | 共享 PostgreSQL 数据库 | 链下索引器写入 events 表 → The Graph 读取 | < 1s |
| 事件索引 GraphQL API | Redis Pub/Sub 实时推送 | 链下索引器处理事件 → 发布到 Redis → API Server 推送到 WebSocket | < 500ms |
| Agent API 服务 | PostgreSQL 查询 | 链下索引器写入 agents/sessions → Agent API 读取 | < 100ms |
| 区块浏览器前端 | GraphQL API | 链下索引器写入 blocks/transactions → GraphQL 查询 | < 50ms |
| 数据分析平台 | TimescaleDB | 链下索引器写入 block_metrics → Grafana 查询 | < 1min |
| S3 冷归档 | 定时任务 | 链下索引器定期归档到 S3 → 恢复时从 S3 读取 | 批量(天级) |
12.3 与 MSG Chain 节点集成的注意事项
- RPC 节点高可用:索引器不应依赖单个 RPC 节点,应使用多节点负载均衡:
┌──────────────────────────────────────────────┐
│ MSG Chain RPC 节点集群 │
│ │
│ Node A (seed) Node B (full) Node C (full) │
│ 10.0.1.1:26657 10.0.1.2:26657 10.0.1.3:26657│
└──────────────────────────────────────────────┘
│ │ │
└──────────────┼──────────────┘
│
Load Balancer
(nginx)
│
▼
indexer (统一入口)
-
连接池配置:索引器到 RPC 节点的连接数应限制,避免压垮节点:
- WebSocket 连接:每个索引器实例 1 个连接
- REST 请求:并发数 ≤ 20
- 重试策略:指数退避 + 抖动
-
节点版本兼容:索引器应验证 MSG Chain 节点版本,兼容 Cosmos SDK v0.47+ 的 REST 和 gRPC 接口
12.4 数据流集成
// src/integration/data-plane-bridge.ts
export class DataPlaneBridge {
private redis: Redis;
private pgPool: Pool;
constructor(redisUrl: string, dbUrl: string) {
this.redis = new Redis(redisUrl);
this.pgPool = new Pool({ connectionString: dbUrl });
}
async publishEventToDataPlane(event: NormalizedEvent): Promise<void> {
const channel = `events:${event.contractAddress}`;
await this.redis.publish(channel, JSON.stringify(event));
await this.redis.xadd(
'events:stream',
'*',
'type', event.eventType,
'contract', event.contractAddress,
'action', event.action,
'data', JSON.stringify(event.attributes),
'height', event.blockHeight.toString(),
);
}
async notifyAgentUpdate(agentId: string, action: string): Promise<void> {
const client = await this.pgPool.connect();
try {
await client.query(
`SELECT pg_notify('agent_updates', $1)`,
[JSON.stringify({ agentId, action, timestamp: new Date().toISOString() })],
);
} finally {
client.release();
}
}
}
13. 总结
13.1 核心结论
| 维度 | 结论 |
|---|---|
| 架构模式 | Active-Passive 为主,分片索引为辅 |
| 部署环境 | K8s (生产) / Docker Compose (开发) |
| Leader 选举 | Redis 分布式锁 + TTL 续约 |
| 存储策略 | PostgreSQL 主存 + Redis 缓存 + S3 冷归档 |
| 数据一致性 | 关键操作强一致,非关键操作最终一致 |
| 故障恢复 | RTO < 60s, RPO < 1 block |
| 监控体系 | Prometheus + Grafana + AlertManager |
| 存储增长 | 预估 ~1.7x 事件数据量/年(含索引和 WAL) |
13.2 推荐演进路径
Phase 1 (最小可行):
- 单索引器 + PostgreSQL + Redis
- 基本监控(健康检查 + 日志)
→ 适用: 开发环境、低流量场景
Phase 2 (高可用):
- Active-Passive 索引器对
- PostgreSQL Primary + Replica
- Prometheus + Grafana 告警
→ 适用: 生产环境、SLA 99.9%
Phase 3 (弹性扩展):
- K8s 部署 + HPA
- 分片索引(按合约拆分)
- Redis 集群 + TimescaleDB
→ 适用: 高吞吐、多合约并行索引
Phase 4 (多活):
- 跨 AZ Active-Passive
- S3 冷归档自动化
- 全局查询加速 (CDN + 边缘缓存)
→ 适用: 全球部署、SLA 99.99%
13.3 参考资源
- MSG Chain 开发者文档: https://msgchain.org/docs
- CosmWasm 文档: https://docs.cosmwasm.com
- The Graph 文档: https://thegraph.com/docs
- PostgreSQL 高可用: https://www.postgresql.org/docs/15/high-availability.html
- Redis Sentinel: https://redis.io/docs/reference/sentinel/
- Kubernetes 部署最佳实践: https://kubernetes.io/docs/concepts/configuration/overview/
本文档由 MSG Chain 开发团队维护。如有问题,请在 GitHub 提交 Issue: https://github.com/msgchain/docs/issues
