Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
File renamed without changes.
3 changes: 1 addition & 2 deletions service/cluster/p2c.go → cluster/p2c.go
Original file line number Diff line number Diff line change
Expand Up @@ -14,7 +14,6 @@ import (
// 代价函数:score = ewmaLatency × (inflight + 1)——既看历史延迟(EWMA 平滑),又看当前在途
// (+1 使空闲后端也有区分度)。冷启动 ewma=0 → score=0 → 优先被选中以探测(类慢启动)。
//
// 与有界负载一致性哈希(BoundedRing)互补:后者管局部性 + 负载上界,本器管延迟感知选优;
// 二者都是无状态请求/副本 LB 原语(如把读请求在一组副本间择优),非有状态数据放置。
//
// 并发安全。
Expand All @@ -23,7 +22,7 @@ type P2CBalancer struct {
backends []string
ewma map[string]float64 // 各后端延迟的 EWMA(纳秒)
inflight map[string]int
decay float64 // EWMA 平滑系数 (0,1],越大越跟新样本
decay float64 // EWMA 平滑系数 (0,1],越大越跟新样本
rng *rand.Rand
}

Expand Down
File renamed without changes.
File renamed without changes.
16 changes: 3 additions & 13 deletions service/cluster/placement.go → cluster/placement.go
Original file line number Diff line number Diff line change
@@ -1,15 +1,11 @@
package cluster

import (
"errors"
"log/slog"
)

// errNotImplemented 标记尚未落地的控制面能力(数据迁移、跨节点转发等)。
// 这些能力依赖传输层重写(当前 Raft 走 net/rpc 静态传输),属 stretch 范围,
// 见架构文档。桩实现统一返回此错误,避免调用方误以为已生效。
var errNotImplemented = errors.New("cluster: not implemented")

// Placement 是放置控制面:组合一致性哈希环与节点注册表,回答「某 key 当前的
// 属主是谁」,并提供故障转移与再平衡的入口。
//
Expand Down Expand Up @@ -45,13 +41,7 @@ func (p *Placement) Failover(deadNode string) {
slog.Info("[cluster] failover: node removed from ring", "node", deadNode)
}

// Rebalance 是再平衡桩(stretch)。
//
// 真正的再平衡需要在成员变更后执行真实的数据迁移(把 key 的实际数据从旧属主
// 搬到新属主),这依赖跨节点数据传输通道——当前 Raft 使用 net/rpc 静态传输,
// 无法承载分片迁移,属传输层重写范围(见架构文档)。此处仅记录 TODO 并返回
// errNotImplemented,保留接口与调用点,待传输层就绪后填充。
func (p *Placement) Rebalance() error {
slog.Warn("[cluster] rebalance: TODO, requires data migration over a new transport (stretch)")
return errNotImplemented
// IsLocal 判定 key 的属主是否为 self(本节点),供网关决定本地处理还是转发。
func (p *Placement) IsLocal(key []byte, self string) bool {
return p.OwnerOf(key) == self
}
File renamed without changes.
File renamed without changes.
22 changes: 16 additions & 6 deletions service/cluster/routing.go → cluster/routing.go
Original file line number Diff line number Diff line change
@@ -1,11 +1,21 @@
// Package cluster 实现分片分布式集群的路由与控制面骨架。
// Package cluster 是分片集群的控制面:决定「一个 key 归属哪个分片、哪个物理节点」,
// 以及节点间的读择优与转发连接复用。它与存储、传输解耦,只依赖 bannet 做跨节点调用。
//
// 设计动机:BanDB 的定位是「数仓写入前置缓冲引擎」,向分片集群演进时需要
// 一层与存储解耦的「放置与路由控制面」(借鉴 dubbo-go / PD 的思路)。本包
// 只承担控制面职责——决定「一个 key 归属哪个物理节点 / 哪个分片」,以及节点
// 存活的注册发现;真实的跨节点数据迁移与传输属于传输层重写范围,本包以桩标注。
// 已在生产路径上运行的部分:
//
// 零第三方依赖:一致性哈希仅使用标准库 hash/crc32。
// - HashRing / ShardOf / ShardReplicas —— 一致性哈希归属与分片副本集,
// 由 service/shardkv 与 service.Router 使用。
// - P(P2C)—— 两选一的延迟感知读择优,由 shardkv 的转发读使用。
// - PeerPool —— 按地址复用的跨节点转发连接池,由 Router 的属主转发使用。
//
// 仍是骨架、当前不产生行为的部分:
//
// - Registry 与 Placement 的存活视图。集群尚无心跳(Heartbeat 无调用方),
// 调用方传入远超进程寿命的 TTL,故所有节点恒被视为存活,Placement.OwnerOf
// 等价于 HashRing.NodeFor。接入心跳后它才开始起作用。
// - Placement.Failover 能把节点摘出环,但目前没有故障检测来触发它。
//
// 零第三方依赖:一致性哈希仅用标准库 hash/crc32。
package cluster

import (
Expand Down
File renamed without changes.
8 changes: 8 additions & 0 deletions docs/distributed-delivery-cluster-skeleton.md
Original file line number Diff line number Diff line change
Expand Up @@ -103,3 +103,11 @@ flowchart TD
## 已知(非本次引入)问题

- `Server/server.go` 无 `//go:build !pprof` 标签,与 `Server/server_pprof.go`(`//go:build pprof`)在 `-tags pprof` 下 `main` 重复声明。这是 origin/main 上的**既有问题**,本次骨架未修(遵循外科手术式修改,单独提出)。修复方式:给 `server.go` 加 `//go:build !pprof`。

---

> **后续状态(本文之后)**:本文所述骨架中,`Placement.Forward` 与 `Placement.Rebalance`
> 两个桩已移除——其理由「跨节点数据传输尚不可用」已不成立,转发能力后来由 `Router` +
> `PeerPool` 经 BanNet 落地(见 iteration-2026-08-05-shard-routing-banNet)。
> `Registry` 与 `Placement.Failover` 保留,但集群仍无心跳,故存活视图当前不产生行为。
> 包位置亦由 `service/cluster` 上提为顶层 `cluster`。
6 changes: 6 additions & 0 deletions docs/iteration-2026-08-05-bounded-load-consistent-hashing.md
Original file line number Diff line number Diff line change
Expand Up @@ -41,3 +41,9 @@ capacity = ⌈(1+ε) · 总负载 / 节点数⌉

- 接入副本读/请求 LB 路径(ShardKV 的读、或转发到副本集)。
- ε 可配、按节点权重的加权有界负载。

---

> **后续状态(本文之后)**:`BoundedRing` 实现已从代码库移除——它自落地起未被任何调用方
> 引用(包括 `cluster` 包内部),实际承担归属计算的一直是 `HashRing`。本文保留为当时的
> 设计与取舍记录;如需重新引入,代码见该次移除前的 git 历史。
101 changes: 0 additions & 101 deletions service/cluster/bounded_ring.go

This file was deleted.

88 changes: 0 additions & 88 deletions service/cluster/bounded_ring_test.go

This file was deleted.

30 changes: 0 additions & 30 deletions service/cluster/gateway.go

This file was deleted.

12 changes: 10 additions & 2 deletions service/cluster_bootstrap.go
Original file line number Diff line number Diff line change
Expand Up @@ -4,10 +4,18 @@ import (
"log/slog"
"time"

"github.com/NeverENG/BanDB/cluster"
"github.com/NeverENG/BanDB/config"
"github.com/NeverENG/BanDB/service/cluster"
)

// assumeAliveTTL 让所有节点恒被视为存活。
//
// 集群目前没有心跳——cluster.Registry.Heartbeat 无任何调用方,故存活视图不会被刷新。
// 此时若取一个有限 TTL,所有节点会在该窗口后被判为失联,Placement.OwnerOf 返回空串,
// 路由随即整体中断。取远超进程寿命的值,是把「尚无心跳」这一事实显式固定下来,
// 而不是伪装成一个会过期的存活窗口。接入心跳后应改为真实的判活窗口。
const assumeAliveTTL = 100 * 365 * 24 * time.Hour

// EnableShardRoutingFromConfig 按配置在 router 上开启分片路由(默认关闭时直接返回)。
// 开启时以 config.Peers 为节点地址构建一致性哈希放置,self = Peers[Me],不属本节点的
// key 经 BanNet 转发到 owner。健康探测/故障转移属后续工作,这里所有节点视为存活。
Expand All @@ -21,7 +29,7 @@ func EnableShardRoutingFromConfig(r *Router) {
return
}
self := peers[config.G.Me]
placement := cluster.NewClusterFromPeers(peers, config.G.VNodes, 100*365*24*time.Hour)
placement := cluster.NewClusterFromPeers(peers, config.G.VNodes, assumeAliveTTL)
pool := cluster.NewPeerPool(5 * time.Second)
r.SetRouting(placement, self, pool)
slog.Info("shard routing enabled", "self", self, "peers", peers)
Expand Down
2 changes: 1 addition & 1 deletion service/router.go
Original file line number Diff line number Diff line change
Expand Up @@ -6,11 +6,11 @@ import (
"log/slog"

"github.com/NeverENG/BanDB/bannet"
"github.com/NeverENG/BanDB/cluster"
"github.com/NeverENG/BanDB/pkg/admission"
"github.com/NeverENG/BanDB/pkg/metrics"
"github.com/NeverENG/BanDB/pkg/predicate"
"github.com/NeverENG/BanDB/pkg/proto"
"github.com/NeverENG/BanDB/service/cluster"
"github.com/NeverENG/BanDB/storage"
)

Expand Down
2 changes: 1 addition & 1 deletion service/shard_routing_integration_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -10,10 +10,10 @@ import (
"time"

"github.com/NeverENG/BanDB/bannet"
"github.com/NeverENG/BanDB/cluster"
"github.com/NeverENG/BanDB/config"
"github.com/NeverENG/BanDB/pkg/predicate"
"github.com/NeverENG/BanDB/pkg/proto"
"github.com/NeverENG/BanDB/service/cluster"
)

// memKV 是隔离的内存 KV,用作每个节点的本地 store——从而在一个进程内起多节点、
Expand Down
2 changes: 1 addition & 1 deletion service/shardkv/read.go
Original file line number Diff line number Diff line change
Expand Up @@ -4,8 +4,8 @@ import (
"fmt"
"time"

"github.com/NeverENG/BanDB/cluster"
"github.com/NeverENG/BanDB/raft"
"github.com/NeverENG/BanDB/service/cluster"
)

// ShardReadArgs / ShardReadReply 是转发读 RPC 的报文:向某分片副本读取一个 key。
Expand Down
2 changes: 1 addition & 1 deletion service/shardkv/shardkv.go
Original file line number Diff line number Diff line change
Expand Up @@ -9,8 +9,8 @@ import (
"sync/atomic"
"time"

"github.com/NeverENG/BanDB/cluster"
"github.com/NeverENG/BanDB/raft"
"github.com/NeverENG/BanDB/service/cluster"
)

// Shard 是一个分片:一个 Raft 组 + 该分片的 FSM store + 一个排空 ApplyCh 的 apply 循环。
Expand Down
Loading