feat(pool): 移除BadgerDB连接池实现并集成NutsDB

- 删除pool/badger.go文件中的BadgerDB连接池相关代码
- 更新WebSocket示例移除badger目录参数依赖
- 更新TCP示例移除badger目录参数依赖
- 升级github.com/dgraph-io/badger/v4依赖从v4.2.0到v4.9.1
- 新增github.com/nutsdb/nutsdb依赖用于替代BadgerDB功能
- 添加WebSocket和TCP连接测试示例代码
- 更新多个间接依赖包版本包括ristretto、humanize、compress等
This commit is contained in:
2026-02-27 10:46:04 +08:00
parent 942eff81fb
commit c50714e8a0
7 changed files with 177 additions and 447 deletions
+46 -1
View File
@@ -18,7 +18,7 @@ func Example() {
}
// 创建TCP服务器
server, err := NewTCPServer("0.0.0.0:8888", config, "./badger/tcp")
server, err := NewTCPServer("0.0.0.0:8888", config)
if err != nil {
fmt.Printf("Failed to create server: %v\n", err)
return
@@ -49,3 +49,48 @@ func Example() {
fmt.Println("TCP server stopped.")
}
// TestTCP 测试TCP连接
func TestTCP() {
fmt.Println("=== 测试TCP连接 ===")
fmt.Println("1. 创建TCP服务器配置")
config := &TcpPoolConfig{
BufferSize: 2048,
MaxConnections: 100000,
ConnectTimeout: time.Second * 5,
ReadTimeout: time.Second * 30,
WriteTimeout: time.Second * 10,
MaxIdleTime: time.Minute * 5,
}
fmt.Println("2. 创建TCP服务器")
server, err := NewTCPServer("0.0.0.0:8888", config)
if err != nil {
fmt.Printf("创建服务器失败:%v\n", err)
return
}
fmt.Println("3. 服务器创建成功")
fmt.Println("4. 获取在线连接数")
count := server.Connection.Count()
fmt.Printf("当前在线连接数:%d\n", count)
fmt.Println("5. 获取所有在线连接ID")
connIDs, err := server.GetAllConnIDs()
if err != nil {
fmt.Printf("获取在线连接ID失败:%v\n", err)
} else {
fmt.Printf("在线连接ID:%v\n", connIDs)
}
fmt.Println("6. 启动服务器")
if err := server.Start(); err != nil {
fmt.Printf("启动服务器失败:%v\n", err)
return
}
fmt.Println("7. 服务器启动成功,运行2秒后停止")
time.Sleep(time.Second * 2)
fmt.Println("8. 停止服务器")
if err := server.Stop(); err != nil {
fmt.Printf("停止服务器失败:%v\n", err)
} else {
fmt.Println("服务器停止成功")
}
fmt.Println("=== TCP测试完成 ===")
}
+27 -27
View File
@@ -33,26 +33,26 @@ type TCPServer struct {
// ConnectionPool 连接池结构
type ConnectionPool struct {
connections map[string]*TcpConnection
badgerPool *pool.BadgerPool
nutsPool *pool.NutsPool
mutex sync.RWMutex
config *TcpPoolConfig
logger *glog.Logger
}
// NewTCPServer 创建一个新的TCP服务器
func NewTCPServer(address string, config *TcpPoolConfig, dbPath string) (*TCPServer, error) {
func NewTCPServer(address string, config *TcpPoolConfig) (*TCPServer, error) {
logger := g.Log(address)
ctx, cancel := context.WithCancel(context.Background())
// 初始化BadgerDB连接池
badgerPool, err := pool.NewBadgerPool(dbPath)
// 初始化NutsDB连接池
nutsPool, err := pool.NewNutsPool()
if err != nil {
return nil, fmt.Errorf("failed to create badger pool: %w", err)
return nil, fmt.Errorf("failed to create nuts pool: %w", err)
}
pool := &ConnectionPool{
connections: make(map[string]*TcpConnection),
badgerPool: badgerPool,
nutsPool: nutsPool,
config: config,
logger: logger,
}
@@ -95,9 +95,9 @@ func (s *TCPServer) Stop() error {
s.Listener.Close()
s.wg.Wait()
s.Connection.Clear()
// 关闭BadgerDB连接池
if err := s.Connection.badgerPool.Close(); err != nil {
s.Logger.Error(s.ctx, fmt.Sprintf("Failed to close BadgerDB pool: %v", err))
// 关闭NutsDB连接池
if err := s.Connection.nutsPool.Close(); err != nil {
s.Logger.Error(s.ctx, fmt.Sprintf("Failed to close NutsDB pool: %v", err))
// 不影响服务器停止,仅记录错误
}
s.Logger.Info(s.ctx, "TCP server stopped")
@@ -123,7 +123,7 @@ func (s *TCPServer) handleConnection(conn *gtcp.Conn) {
s.Connection.Add(tcpConn)
s.Logger.Info(s.ctx, fmt.Sprintf("New connection established: %s", connID))
// 存储到BadgerDB
// 存储到NutsDB
connInfo := &pool.ConnectionInfo{
ID: connID,
Type: pool.ConnTypeTCP,
@@ -135,8 +135,8 @@ func (s *TCPServer) handleConnection(conn *gtcp.Conn) {
"localAddress": conn.LocalAddr().String(),
},
}
if err := s.Connection.badgerPool.Add(connInfo); err != nil {
s.Logger.Error(s.ctx, fmt.Sprintf("Failed to store connection to BadgerDB: %v", err))
if err := s.Connection.nutsPool.Add(connInfo); err != nil {
s.Logger.Error(s.ctx, fmt.Sprintf("Failed to store connection to NutsDB: %v", err))
// 不影响连接建立,仅记录错误
}
@@ -152,9 +152,9 @@ func (s *TCPServer) receiveMessages(conn *TcpConnection) {
}
s.Connection.Remove(conn.Id)
conn.Server.Close()
// 从BadgerDB移除
if err := s.Connection.badgerPool.Remove(conn.Id); err != nil {
s.Logger.Error(s.ctx, fmt.Sprintf("Failed to remove connection from BadgerDB: %v", err))
// 从NutsDB移除
if err := s.Connection.nutsPool.Remove(conn.Id); err != nil {
s.Logger.Error(s.ctx, fmt.Sprintf("Failed to remove connection from NutsDB: %v", err))
// 不影响连接关闭,仅记录错误
}
s.Logger.Info(s.ctx, fmt.Sprintf("Connection closed: %s", conn.Id))
@@ -183,12 +183,12 @@ func (s *TCPServer) receiveMessages(conn *TcpConnection) {
conn.LastUsed = now
conn.Mutex.Unlock()
// 更新BadgerDB中的连接信息
connInfo, err := s.Connection.badgerPool.Get(conn.Id)
// 更新NutsDB中的连接信息
connInfo, err := s.Connection.nutsPool.Get(conn.Id)
if err == nil && connInfo != nil {
connInfo.LastUsed = now
if err := s.Connection.badgerPool.Update(connInfo); err != nil {
s.Logger.Error(s.ctx, fmt.Sprintf("Failed to update connection in BadgerDB: %v", err))
if err := s.Connection.nutsPool.Update(connInfo); err != nil {
s.Logger.Error(s.ctx, fmt.Sprintf("Failed to update connection in NutsDB: %v", err))
// 不影响消息处理,仅记录错误
}
}
@@ -259,12 +259,12 @@ func (s *TCPServer) sendMessage(conn *TcpConnection, data []byte) error {
now := time.Now()
conn.LastUsed = now
// 更新BadgerDB中的连接信息
connInfo, err := s.Connection.badgerPool.Get(conn.Id)
// 更新NutsDB中的连接信息
connInfo, err := s.Connection.nutsPool.Get(conn.Id)
if err == nil && connInfo != nil {
connInfo.LastUsed = now
if err := s.Connection.badgerPool.Update(connInfo); err != nil {
s.Logger.Error(s.ctx, fmt.Sprintf("Failed to update connection in BadgerDB: %v", err))
if err := s.Connection.nutsPool.Update(connInfo); err != nil {
s.Logger.Error(s.ctx, fmt.Sprintf("Failed to update connection in NutsDB: %v", err))
// 不影响消息发送,仅记录错误
}
}
@@ -283,9 +283,9 @@ func (s *TCPServer) Kick(connID string) error {
conn.Server.Close()
// 从连接池移除
s.Connection.Remove(connID)
// 从BadgerDB移除
if err := s.Connection.badgerPool.Remove(connID); err != nil {
s.Logger.Error(s.ctx, fmt.Sprintf("Failed to remove connection from BadgerDB: %v", err))
// 从NutsDB移除
if err := s.Connection.nutsPool.Remove(connID); err != nil {
s.Logger.Error(s.ctx, fmt.Sprintf("Failed to remove connection from NutsDB: %v", err))
// 不影响连接关闭,仅记录错误
}
@@ -350,5 +350,5 @@ func (p *ConnectionPool) Count() int {
// GetAllConnIDs 获取所有在线连接的ID列表
func (p *ConnectionPool) GetAllConnIDs() ([]string, error) {
return p.badgerPool.GetAllConnIDs()
return p.nutsPool.GetAllConnIDs()
}