Leader Lease + Metadata Quorum + 2PC
leader lease 决定每个分片的写入归属,metadata quorum 记录事务投票, 2PC 推进跨分片提交,decision replay 负责收尾及故障后的状态对齐。
- 01. lease 决定分片的写入归属
- 02. 2PC 保证跨分片 prepare / commit 的原子性
- 03. replay 在节点重启或局部故障后补齐事务状态
distributed transactions
控制面负责协调、提交与重放。相关逻辑在 benchmark、故障演练与恢复流程中被持续执行,而非仅作为静态代码存在。
leader lease 决定每个分片的写入归属,metadata quorum 记录事务投票, 2PC 推进跨分片提交,decision replay 负责收尾及故障后的状态对齐。
该机制的实际效果:分布式写入可协调、可重放、可恢复,避免多节点各自写入导致的状态不一致。
STATE 01 PREPARE -> lease owner validates shard intent
STATE 02 QUORUM -> metadata quorum records transaction vote
STATE 03 COMMIT -> two-phase commit advances all participants
STATE 04 REPLAY -> recovering nodes reconcile final decision
$ tsdb-admin txn verify --mode strict --shard 0x09ff2
online recovery
恢复逻辑与写路径、控制面共用同一套代码,纳入日常回归。节点宕机重启后,系统通过 repair 机制补齐副本数据。
通过 `/cluster/recovery/state` 查询各分片的修复进度。
恢复能力由自动化故障演练验证:终止节点进程,校验恢复结果。
[t29] online recovery convergence = 34ms
[t30] repair + state stabilization = 2072ms
[state] periodic leader repair = enabled
[api] /cluster/recovery/state = observable
副本修复由 leader 周期性触发,恢复进度通过状态接口暴露, t29 / t30 两组耗时指标随每次回归更新。
stream + sdk
流式订阅与 SDK 是外部使用该系统的主要入口,均配有独立的功能测试。
流写入依赖 sequence 与 replay buffer:连接断开重连后,缺失数据会重新推送; 当请求到达非 owner 节点时,会被转发至正确节点。提供 Python、JDBC、REST 三种客户端,并支持 Rust crate 与 C / C++ 协议级接入。
from tsdb import TsdbClient, read_db
client = TsdbClient("http://127.0.0.1:8080")
client.health()
frame = read_db(
"SELECT * FROM trades LIMIT 10",
url="http://127.0.0.1:8080",
)
jdbc:tsdb://127.0.0.1:8080
以上为 Python 与 JDBC 的最小接入示例,更多用法见仓库 python/ 与 jdbc/ 目录。
benchmark
以上为 2026-08-25 冻结构建的同机对打结果(432 万行,rows_match 三口径门禁全绿;1206 单元测试通过)。完整报告见仓库 bench/ 与 docs/ 目录。
usage
cargo build --release --bin tsdb-server
./target/release/tsdb-server \
--addr 127.0.0.1:8080 \
--data-dir ./data
from tsdb import TsdbClient, read_db
client = TsdbClient("http://127.0.0.1:8080")
print(client.health())
frame = read_db(
"SELECT * FROM trades LIMIT 10",
url="http://127.0.0.1:8080",
)
Connection conn =
DriverManager.getConnection(
"jdbc:tsdb://127.0.0.1:8080");
Statement stmt = conn.createStatement();
ResultSet rs = stmt.executeQuery(
"SELECT * FROM trades LIMIT 10");
use tsdb::network::client::Client;
#[tokio::main]
async fn main() -> tsdb::Result<()> {
let client =
Client::new("http://127.0.0.1:8080");
client.health_check().await?;
let result = client
.execute_query(
"SELECT * FROM trades LIMIT 10")
.await?;
Ok(())
}
/* 与 tsdb_fdw 相同的方式:
libcurl POST /jdbc */
CURL *curl = curl_easy_init();
struct curl_slist *h = NULL;
h = curl_slist_append(h,
"Content-Type: application/json");
curl_easy_setopt(curl, CURLOPT_URL,
"http://127.0.0.1:8080/jdbc");
curl_easy_setopt(curl,
CURLOPT_HTTPHEADER, h);
curl_easy_setopt(curl, CURLOPT_POSTFIELDS,
"{\"action\":\"execute\","
"\"sql\":\"SELECT * FROM trades"
" LIMIT 10\"}");
curl_easy_perform(curl);
// REST 接入;大批量写入可走
// TCP 二进制协议 :9090
CURL *curl = curl_easy_init();
curl_easy_setopt(curl, CURLOPT_URL,
"http://127.0.0.1:8080/jdbc");
std::string body = R"({
"action": "execute",
"sql": "SELECT * FROM trades LIMIT 10"
})";
curl_easy_setopt(curl,
CURLOPT_POSTFIELDS, body.c_str());
curl_easy_perform(curl);