Go语言实现高并发TCP连接池:连接复用与生命周期管理

为什么需要TCP连接池

每次建立TCP连接需要经历三次握手,在高并发场景下,频繁创建和销毁连接的开销显著。连接池通过复用已建立的连接,消除了握手延迟和端口消耗问题。Go语言的标准库没有提供通用的TCP连接池,需要自行实现或使用第三方库。本文从零实现一个支持连接复用、健康检查和优雅关闭的TCP连接池。

连接池核心结构设计

package pool

import (
    "errors"
    "net"
    "sync"
    "time"
)

var (
    ErrPoolClosed    = errors.New("pool is closed")
    ErrPoolExhausted = errors.New("pool exhausted")
)

type PoolConn struct {
    net.Conn
    pool       *TCPPool
    createdAt  time.Time
    lastUsedAt time.Time
}

func (pc *PoolConn) Close() error {
    return pc.pool.put(pc)
}

type Config struct {
    Addr                string
    InitialCap          int
    MaxCap              int
    MaxIdle             int
    IdleTimeout         time.Duration
    ConnectTimeout      time.Duration
    HealthCheckInterval time.Duration
}

type TCPPool struct {
    config  Config
    mu      sync.Mutex
    conns   chan *PoolConn
    active  int
    closed  bool
    factory func() (net.Conn, error)
    stopCh  chan struct{}
}

func NewTCPPool(config Config) (*TCPPool, error) {
    if config.InitialCap < 0 || config.MaxCap <= 0 {
        return nil, errors.New("invalid capacity config")
    }
    if config.InitialCap > config.MaxCap {
        return nil, errors.New("initial capacity exceeds max")
    }

    p := &TCPPool{
        config:  config,
        conns:   make(chan *PoolConn, config.MaxCap),
        factory: func() (net.Conn, error) {
            return net.DialTimeout("tcp", config.Addr, config.ConnectTimeout)
        },
        stopCh: make(chan struct{}),
    }

    for i := 0; i < config.InitialCap; i++ {
        conn, err := p.factory()
        if err != nil {
            p.Close()
            return nil, err
        }
        p.conns <- &PoolConn{
            Conn:       conn,
            pool:       p,
            createdAt:  time.Now(),
            lastUsedAt: time.Now(),
        }
        p.active++
    }

    if config.HealthCheckInterval > 0 {
        go p.healthCheck()
    }

    return p, nil
}

获取与归还连接

Get方法从池中获取连接,如果没有空闲连接且未达到上限,则新建一个:

func (p *TCPPool) Get() (*PoolConn, error) {
    p.mu.Lock()
    if p.closed {
        p.mu.Unlock()
        return nil, ErrPoolClosed
    }
    p.mu.Unlock()

    // 尝试从缓冲通道获取空闲连接
    select {
    case pc := <-p.conns:
        if p.isConnAlive(pc) {
            pc.lastUsedAt = time.Now()
            return pc, nil
        }
        pc.Conn.Close()
        p.mu.Lock()
        p.active--
        p.mu.Unlock()
    default:
    }

    // 没有空闲连接,尝试新建
    p.mu.Lock()
    if p.active >= p.config.MaxCap {
        p.mu.Unlock()
        select {
        case pc := <-p.conns:
            pc.lastUsedAt = time.Now()
            return pc, nil
        case <-time.After(5 * time.Second):
            return nil, ErrPoolExhausted
        }
    }
    p.active++
    p.mu.Unlock()

    conn, err := p.factory()
    if err != nil {
        p.mu.Lock()
        p.active--
        p.mu.Unlock()
        return nil, err
    }

    return &PoolConn{
        Conn:       conn,
        pool:       p,
        createdAt:  time.Now(),
        lastUsedAt: time.Now(),
    }, nil
}

func (p *TCPPool) put(pc *PoolConn) error {
    p.mu.Lock()
    if p.closed {
        p.mu.Unlock()
        pc.Conn.Close()
        p.mu.Lock()
        p.active--
        p.mu.Unlock()
        return ErrPoolClosed
    }
    p.mu.Unlock()

    if !p.isConnAlive(pc) {
        pc.Conn.Close()
        p.mu.Lock()
        p.active--
        p.mu.Unlock()
        return nil
    }

    select {
    case p.conns <- pc:
        return nil
    default:
        pc.Conn.Close()
        p.mu.Lock()
        p.active--
        p.mu.Unlock()
        return nil
    }
}

连接健康检查与空闲回收

长时间空闲的连接可能已被防火墙或对端关闭,继续使用会报错。定期健康检查可以提前发现并清理这些半开连接:

func (p *TCPPool) healthCheck() {
    ticker := time.NewTicker(p.config.HealthCheckInterval)
    defer ticker.Stop()

    for {
        select {
        case <-ticker.C:
            p.cleanIdleConns()
        case <-p.stopCh:
            return
        }
    }
}

func (p *TCPPool) cleanIdleConns() {
    p.mu.Lock()
    defer p.mu.Unlock()

    if p.closed {
        return
    }

    var keep []*PoolConn
    for i := 0; i < len(p.conns); i++ {
        select {
        case pc := <-p.conns:
            if p.config.IdleTimeout > 0 &&
               time.Since(pc.lastUsedAt) > p.config.IdleTimeout {
                pc.Conn.Close()
                p.active--
                continue
            }
            if !p.isConnAlive(pc) {
                pc.Conn.Close()
                p.active--
                continue
            }
            keep = append(keep, pc)
        default:
            break
        }
    }

    for _, pc := range keep {
        p.conns <- pc
    }
}

优雅关闭连接池

服务关闭时需要等待正在使用的连接归还,然后统一关闭所有连接:

func (p *TCPPool) Close() error {
    p.mu.Lock()
    if p.closed {
        p.mu.Unlock()
        return nil
    }
    p.closed = true
    close(p.stopCh)
    p.mu.Unlock()

    close(p.conns)
    var errs []error
    for pc := range p.conns {
        if err := pc.Conn.Close(); err != nil {
            errs = append(errs, err)
        }
        p.mu.Lock()
        p.active--
        p.mu.Unlock()
    }

    if len(errs) > 0 {
        return errs[0]
    }
    return nil
}

连接池的压测与调优

连接池的参数调优依赖实际压测数据。关键参数包括MaxCap(最大连接数)、IdleTimeout(空闲超时)和InitialCap(预热连接数):

# 使用 vegeta 压测
# echo "GET http://localhost:8080/api" | vegeta attack -duration=30s -rate=5000 | vegeta report
#
# 连接池参数推荐基线:
# MaxCap = CPU核心数 * 50(Go的调度器能有效利用此数量)
# IdleTimeout = 30s(短于对端防火墙的空闲超时)
# InitialCap = MaxCap / 4(避免启动时大量并发建连)

通过监控连接池的active和idle数量曲线,可以判断MaxCap是否足够(active持续等于MaxCap说明需要扩容),IdleTimeout是否合理(idle数量在低峰期过多说明IdleTimeout设置过长或MaxIdle过大)。

原创文章,作者:小编,如若转载,请注明出处:https://www.yunthe.com/go-yu-yan-shi-xian-gao-bing-fa-tcp-lian-jie-chi-lian-jie-fu/

(0)
小编小编
上一篇 1天前
下一篇 1天前

相关推荐