Mysql 几种主从复制方式

一、MySQL主从复制方式

1. 异步复制(默认方式)

  • 原理:主库执行完事务后立即返回客户端,不等待从库接收binlog

  • 优点:性能好,对主库性能影响小

  • 缺点数据可能丢失(主库宕机时未同步的数据)

  • 配置示例

# 主库配置(my.cnf)
server-id=1
log-bin=mysql-bin
binlog-format=ROW

# 从库配置
server-id=2
relay-log=mysql-relay-bin
read-only=1

2. 半同步复制

  • 原理:主库提交事务后至少等待一个从库接收binlog才返回客户端

  • 优点:数据安全性更高

  • 缺点:性能略有下降,网络延迟影响响应时间

  • 启用命令

# 安装插件(主从都需要)
INSTALL PLUGIN rpl_semi_sync_master SONAME 'semisync_master.so';
INSTALL PLUGIN rpl_semi_sync_slave SONAME 'semisync_slave.so';

# 开启半同步
SET GLOBAL rpl_semi_sync_master_enabled=1;
SET GLOBAL rpl_semi_sync_slave_enabled=1;

3. 同步复制(MySQL Cluster/NDB)

  • 原理:所有节点同步完成后才返回客户端

  • 场景:金融级强一致性要求

4. GTID复制(推荐)

  • 原理:使用全局事务ID标识每个事务,简化复制管理和故障恢复

  • 配置

# 主从库都配置
gtid_mode=ON
enforce_gtid_consistency=ON

# 从库配置
CHANGE MASTER TO
MASTER_HOST='master_ip',
MASTER_USER='repl',
MASTER_PASSWORD='password',
MASTER_AUTO_POSITION=1;

5. 基于位点的复制(传统方式)

CHANGE MASTER TO
MASTER_HOST='master_ip',
MASTER_LOG_FILE='mysql-bin.000001',
MASTER_LOG_POS=107;

二、Go Gin框架实践案例

项目结构

project/
├── main.go
├── config/
├── middleware/
├── model/
├── router/
└── pkg/
    └── database/

1. 数据库连接池配置

// pkg/database/database.go
package database

import (
    "gorm.io/driver/mysql"
    "gorm.io/gorm"
    "gorm.io/gorm/logger"
)

type DBConfig struct {
    Master string   `yaml:"master"` // 主库DSN
    Slaves []string `yaml:"slaves"` // 从库DSN列表
}

var (
    MasterDB *gorm.DB
    SlaveDBs []*gorm.DB
    slaveIndex = 0
)

// 初始化数据库连接
func InitDB(config DBConfig) error {
    // 连接主库(写操作)
    masterDB, err := gorm.Open(mysql.Open(config.Master), &gorm.Config{
        Logger: logger.Default.LogMode(logger.Info),
    })
    if err != nil {
        return err
    }
    MasterDB = masterDB
    
    // 连接从库(读操作)
    for _, slaveDSN := range config.Slaves {
        slaveDB, err := gorm.Open(mysql.Open(slaveDSN), &gorm.Config{
            Logger: logger.Default.LogMode(logger.Info),
        })
        if err != nil {
            return err
        }
        SlaveDBs = append(SlaveDBs, slaveDB)
    }
    
    return nil
}

// 获取从库(轮询负载均衡)
func GetSlaveDB() *gorm.DB {
    if len(SlaveDBs) == 0 {
        return MasterDB
    }
    
    db := SlaveDBs[slaveIndex]
    slaveIndex = (slaveIndex + 1) % len(SlaveDBs)
    return db
}

2. Gin中间件实现读写分离

// middleware/db_router.go
package middleware

import (
    "github.com/gin-gonic/gin"
    "net/http"
    "strings"
    "your-project/pkg/database"
)

// DBContextKey 数据库上下文键
type DBContextKey string

const (
    WriteDB DBContextKey = "write_db"
    ReadDB  DBContextKey = "read_db"
)

// DBRouterMiddleware 数据库路由中间件
func DBRouterMiddleware() gin.HandlerFunc {
    return func(c *gin.Context) {
        // 根据HTTP方法决定使用主库还是从库
        method := c.Request.Method
        
        // 写操作使用主库
        if method == http.MethodPost || 
           method == http.MethodPut || 
           method == http.MethodPatch || 
           method == http.MethodDelete {
            c.Set(string(WriteDB), database.MasterDB)
            c.Set(string(ReadDB), database.MasterDB)
        } else {
            // 读操作使用从库
            c.Set(string(WriteDB), database.MasterDB)
            c.Set(string(ReadDB), database.GetSlaveDB())
        }
        
        c.Next()
    }
}

// 获取数据库连接的辅助函数
func GetWriteDB(c *gin.Context) *gorm.DB {
    if db, exists := c.Get(string(WriteDB)); exists {
        return db.(*gorm.DB)
    }
    return database.MasterDB
}

func GetReadDB(c *gin.Context) *gorm.DB {
    if db, exists := c.Get(string(ReadDB)); exists {
        return db.(*gorm.DB)
    }
    return database.GetSlaveDB()
}

3. 模型定义与使用

// model/user.go
package model

import (
    "github.com/gin-gonic/gin"
    "gorm.io/gorm"
    "your-project/middleware"
)

type User struct {
    ID       uint   `gorm:"primaryKey" json:"id"`
    Username string `gorm:"size:50" json:"username"`
    Email    string `gorm:"size:100;unique" json:"email"`
}

// 创建用户(写操作)
func CreateUser(c *gin.Context, user *User) error {
    db := middleware.GetWriteDB(c)
    return db.Create(user).Error
}

// 获取用户列表(读操作)
func GetUsers(c *gin.Context, page, pageSize int) ([]User, int64, error) {
    var users []User
    var total int64
    
    db := middleware.GetReadDB(c)
    
    // 获取总数
    if err := db.Model(&User{}).Count(&total).Error; err != nil {
        return nil, 0, err
    }
    
    // 分页查询
    offset := (page - 1) * pageSize
    err := db.Offset(offset).Limit(pageSize).Find(&users).Error
    
    return users, total, err
}

// 获取单个用户(读操作)
func GetUserByID(c *gin.Context, id uint) (*User, error) {
    var user User
    db := middleware.GetReadDB(c)
    err := db.First(&user, id).Error
    return &user, err
}

4. 路由控制器

// controller/user_controller.go
package controller

import (
    "github.com/gin-gonic/gin"
    "net/http"
    "strconv"
    "your-project/model"
)

type UserController struct{}

func (uc *UserController) RegisterRoutes(r *gin.RouterGroup) {
    r.POST("/users", uc.CreateUser)
    r.GET("/users", uc.GetUserList)
    r.GET("/users/:id", uc.GetUser)
}

func (uc *UserController) CreateUser(c *gin.Context) {
    var user model.User
    if err := c.ShouldBindJSON(&user); err != nil {
        c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()})
        return
    }
    
    if err := model.CreateUser(c, &user); err != nil {
        c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
        return
    }
    
    c.JSON(http.StatusCreated, user)
}

func (uc *UserController) GetUserList(c *gin.Context) {
    page, _ := strconv.Atoi(c.DefaultQuery("page", "1"))
    pageSize, _ := strconv.Atoi(c.DefaultQuery("page_size", "10"))
    
    users, total, err := model.GetUsers(c, page, pageSize)
    if err != nil {
        c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
        return
    }
    
    c.JSON(http.StatusOK, gin.H{
        "data":  users,
        "total": total,
        "page":  page,
    })
}

func (uc *UserController) GetUser(c *gin.Context) {
    id, err := strconv.ParseUint(c.Param("id"), 10, 32)
    if err != nil {
        c.JSON(http.StatusBadRequest, gin.H{"error": "invalid id"})
        return
    }
    
    user, err := model.GetUserByID(c, uint(id))
    if err != nil {
        c.JSON(http.StatusNotFound, gin.H{"error": "user not found"})
        return
    }
    
    c.JSON(http.StatusOK, user)
}

5. 主程序入口

// main.go
package main

import (
    "github.com/gin-gonic/gin"
    "gopkg.in/yaml.v3"
    "io/ioutil"
    "log"
    "your-project/config"
    "your-project/controller"
    "your-project/middleware"
    "your-project/pkg/database"
    "your-project/router"
)

func main() {
    // 读取配置
    cfg := loadConfig()
    
    // 初始化数据库
    dbConfig := database.DBConfig{
        Master: cfg.Database.Master,
        Slaves: cfg.Database.Slaves,
    }
    
    if err := database.InitDB(dbConfig); err != nil {
        log.Fatal("Failed to connect database:", err)
    }
    
    // 创建Gin应用
    app := gin.Default()
    
    // 使用数据库路由中间件
    app.Use(middleware.DBRouterMiddleware())
    
    // 注册路由
    api := app.Group("/api")
    {
        userCtrl := &controller.UserController{}
        userCtrl.RegisterRoutes(api.Group("/v1"))
    }
    
    // 启动服务
    log.Println("Server starting on :8080")
    app.Run(":8080")
}

func loadConfig() *config.Config {
    data, err := ioutil.ReadFile("config.yaml")
    if err != nil {
        log.Fatal("Failed to read config:", err)
    }
    
    var cfg config.Config
    if err := yaml.Unmarshal(data, &cfg); err != nil {
        log.Fatal("Failed to parse config:", err)
    }
    
    return &cfg
}

6. 配置文件示例

# config.yaml
database:
  master: "username:password@tcp(127.0.0.1:3306)/mydb?charset=utf8mb4&parseTime=True&loc=Local"
  slaves:
    - "username:password@tcp(127.0.0.1:3307)/mydb?charset=utf8mb4&parseTime=True&loc=Local"
    - "username:password@tcp(127.0.0.1:3308)/mydb?charset=utf8mb4&parseTime=True&loc=Local"

server:
  port: 8080
  mode: "debug"

三、注意事项

1. 数据一致性考虑

// 强制读主库的场景(如刚写入后立即读取)
func GetUserAfterUpdate(c *gin.Context, id uint) (*User, error) {
    // 某些业务场景需要读主库保证一致性
    db := middleware.GetWriteDB(c)
    var user User
    err := db.First(&user, id).Error
    return &user, err
}

2. 连接池配置优化

sqlDB, err := MasterDB.DB()
if err == nil {
    // 设置连接池参数
    sqlDB.SetMaxIdleConns(10)
    sqlDB.SetMaxOpenConns(100)
    sqlDB.SetConnMaxLifetime(time.Hour)
}

3. 监控与健康检查

// 添加从库健康检查
func CheckSlaveHealth() {
    for i, slave := range database.SlaveDBs {
        if err := slave.Exec("SELECT 1").Error; err != nil {
            log.Printf("Slave %d is down: %v", i, err)
            // 可以从SlaveDBs中移除故障从库
        }
    }
}

四、最佳实践建议

  • 读写分离策略

    • 读多写少的场景适合主从架构

    • 强一致性要求高的业务慎用从库读

  • 延迟处理

    • 主从同步有延迟,业务需要容忍

    • 关键业务操作后立即读取应走主库

  • 故障切换

    • 实现从库健康检查

    • 主库故障时要有应急预案

  • 监控指标

    • 监控主从延迟(Seconds_Behind_Master)

    • 监控数据库连接数

    • 监控查询性能

这种架构能够有效分摊数据库压力,提高系统吞吐量,特别适合读多写少的Web应用场景。

使用Docker搭建MySQL主从复制环境,并在Gin框架中实践读写分离

一、Docker搭建MySQL主从复制环境

1. 创建Docker网络

docker network create mysql-replication-net

2. 创建主从配置文件

主库配置 (master.cnf)

# master.cnf
[mysqld]
server-id=1
log-bin=mysql-bin
binlog-format=ROW
expire_logs_days=7
max_binlog_size=100M
binlog_do_db=app_db
sync_binlog=1
innodb_flush_log_at_trx_commit=1
character-set-server=utf8mb4
collation-server=utf8mb4_unicode_ci

[client]
default-character-set=utf8mb4

从库配置 (slave.cnf

# slave.cnf
[mysqld]
server-id=2
relay-log=mysql-relay-bin
read-only=1
log_slave_updates=1
skip_slave_start=0
character-set-server=utf8mb4
collation-server=utf8mb4_unicode_ci
max_connections=1000
innodb_buffer_pool_size=256M

[client]
default-character-set=utf8mb4

3. 创建docker-compose.yml

# docker-compose.yml
version: '3.8'

services:
  mysql-master:
    image: mysql:8.0
    container_name: mysql-master
    restart: always
    environment:
      MYSQL_ROOT_PASSWORD: root123
      MYSQL_DATABASE: app_db
      MYSQL_USER: app_user
      MYSQL_PASSWORD: app_pass123
    ports:
      - "3307:3306"
    volumes:
      - ./master.cnf:/etc/mysql/conf.d/master.cnf:ro
      - mysql-master-data:/var/lib/mysql
      - ./init:/docker-entrypoint-initdb.d
    networks:
      - mysql-replication-net
    command: 
      - --default-authentication-plugin=mysql_native_password
      - --character-set-server=utf8mb4
      - --collation-server=utf8mb4_unicode_ci
    healthcheck:
      test: ["CMD", "mysqladmin", "ping", "-h", "localhost", "-uroot", "-proot123"]
      interval: 10s
      timeout: 5s
      retries: 5

  mysql-slave1:
    image: mysql:8.0
    container_name: mysql-slave1
    restart: always
    environment:
      MYSQL_ROOT_PASSWORD: root123
      MYSQL_DATABASE: app_db
    ports:
      - "3308:3306"
    volumes:
      - ./slave.cnf:/etc/mysql/conf.d/slave.cnf:ro
      - mysql-slave1-data:/var/lib/mysql
    networks:
      - mysql-replication-net
    command: 
      - --default-authentication-plugin=mysql_native_password
      - --character-set-server=utf8mb4
      - --collation-server=utf8mb4_unicode_ci
    depends_on:
      mysql-master:
        condition: service_healthy
    healthcheck:
      test: ["CMD", "mysqladmin", "ping", "-h", "localhost", "-uroot", "-proot123"]
      interval: 10s
      timeout: 5s
      retries: 5

  mysql-slave2:
    image: mysql:8.0
    container_name: mysql-slave2
    restart: always
    environment:
      MYSQL_ROOT_PASSWORD: root123
      MYSQL_DATABASE: app_db
    ports:
      - "3309:3306"
    volumes:
      - ./slave.cnf:/etc/mysql/conf.d/slave.cnf:ro
      - mysql-slave2-data:/var/lib/mysql
    networks:
      - mysql-replication-net
    command: 
      - --default-authentication-plugin=mysql_native_password
      - --character-set-server=utf8mb4
      - --collation-server=utf8mb4_unicode_ci
    depends_on:
      mysql-master:
        condition: service_healthy
    healthcheck:
      test: ["CMD", "mysqladmin", "ping", "-h", "localhost", "-uroot", "-proot123"]
      interval: 10s
      timeout: 5s
      retries: 5

  mysql-admin:
    image: phpmyadmin/phpmyadmin
    container_name: mysql-admin
    restart: always
    ports:
      - "8081:80"
    environment:
      PMA_HOST: mysql-master
      PMA_PORT: 3306
      UPLOAD_LIMIT: 100M
    networks:
      - mysql-replication-net
    depends_on:
      - mysql-master

networks:
  mysql-replication-net:
    external: true

volumes:
  mysql-master-data:
  mysql-slave1-data:
  mysql-slave2-data:

4. 创建初始化脚本

# init/01-init-master.sql
-- 创建复制用户
CREATE USER 'repl_user'@'%' IDENTIFIED BY 'repl_pass123';
GRANT REPLICATION SLAVE ON *.* TO 'repl_user'@'%';
FLUSH PRIVILEGES;

-- 创建应用用户
GRANT ALL PRIVILEGES ON app_db.* TO 'app_user'@'%';
FLUSH PRIVILEGES;

-- 创建测试表
CREATE DATABASE IF NOT EXISTS app_db;
USE app_db;

CREATE TABLE users (
    id BIGINT AUTO_INCREMENT PRIMARY KEY,
    username VARCHAR(50) NOT NULL UNIQUE,
    email VARCHAR(100) NOT NULL UNIQUE,
    created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
    updated_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;

CREATE TABLE products (
    id BIGINT AUTO_INCREMENT PRIMARY KEY,
    name VARCHAR(100) NOT NULL,
    price DECIMAL(10, 2) NOT NULL,
    stock INT DEFAULT 0,
    created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
    updated_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;

5. 创建从库设置脚本

# setup-replication.sh
#!/bin/bash

echo "=== 设置MySQL主从复制 ==="

# 获取主库状态
MASTER_STATUS=$(docker exec mysql-master mysql -uroot -proot123 -e "SHOW MASTER STATUS\G")

echo "主库状态:"
echo "$MASTER_STATUS"

# 提取File和Position
MASTER_LOG_FILE=$(echo "$MASTER_STATUS" | grep "File:" | awk '{print $2}')
MASTER_LOG_POS=$(echo "$MASTER_STATUS" | grep "Position:" | awk '{print $2}')

echo "Master Log File: $MASTER_LOG_FILE"
echo "Master Log Position: $MASTER_LOG_POS"

# 设置从库1
echo "设置从库1..."
docker exec mysql-slave1 mysql -uroot -proot123 <<EOF
STOP SLAVE;
CHANGE MASTER TO
MASTER_HOST='mysql-master',
MASTER_USER='repl_user',
MASTER_PASSWORD='repl_pass123',
MASTER_LOG_FILE='$MASTER_LOG_FILE',
MASTER_LOG_POS=$MASTER_LOG_POS;
START SLAVE;
EOF

echo "从库1状态:"
docker exec mysql-slave1 mysql -uroot -proot123 -e "SHOW SLAVE STATUS\G" | grep -E "Slave_IO_Running|Slave_SQL_Running"

# 设置从库2
echo "设置从库2..."
docker exec mysql-slave2 mysql -uroot -proot123 <<EOF
STOP SLAVE;
CHANGE MASTER TO
MASTER_HOST='mysql-master',
MASTER_USER='repl_user',
MASTER_PASSWORD='repl_pass123',
MASTER_LOG_FILE='$MASTER_LOG_FILE',
MASTER_LOG_POS=$MASTER_LOG_POS;
START SLAVE;
EOF

echo "从库2状态:"
docker exec mysql-slave2 mysql -uroot -proot123 -e "SHOW SLAVE STATUS\G" | grep -E "Slave_IO_Running|Slave_SQL_Running"

echo "=== 设置完成 ==="

6. 启动和配置

# 1. 启动服务
docker-compose up -d

# 2. 等待MySQL启动
sleep 30

# 3. 设置复制权限
chmod +x setup-replication.sh
./setup-replication.sh

# 4. 验证复制状态
docker exec mysql-slave1 mysql -uroot -proot123 -e "SHOW SLAVE STATUS\G" | grep -A5 -B5 "Running"

二、Go Gin项目完整实现

项目结构

gin-mysql-replication/
├── cmd/
│   └── server/
│       └── main.go
├── config/
│   ├── config.go
│   └── config.yaml
├── internal/
│   ├── middleware/
│   │   ├── db_router.go
│   │   └── recovery.go
│   ├── model/
│   │   ├── user.go
│   │   ├── product.go
│   │   └── base.go
│   ├── repository/
│   │   ├── user_repo.go
│   │   ├── product_repo.go
│   │   └── base_repo.go
│   ├── service/
│   │   ├── user_service.go
│   │   ├── product_service.go
│   │   └── health_service.go
│   ├── handler/
│   │   ├── user_handler.go
│   │   ├── product_handler.go
│   │   └── health_handler.go
│   └── router/
│       └── router.go
├── pkg/
│   ├── database/
│   │   └── db.go
│   └── logger/
│       └── logger.go
├── docker-compose.yml
├── master.cnf
├── slave.cnf
├── setup-replication.sh
├── Dockerfile
├── go.mod
└── README.md

1. 数据库连接池封装

// pkg/database/db.go
package database

import (
    "context"
    "fmt"
    "sync"
    "time"

    "gorm.io/driver/mysql"
    "gorm.io/gorm"
    "gorm.io/gorm/logger"
    "gorm.io/gorm/schema"
)

type DBConfig struct {
    Master      string   `yaml:"master"`
    Slaves      []string `yaml:"slaves"`
    MaxIdleConn int      `yaml:"max_idle_conn"`
    MaxOpenConn int      `yaml:"max_open_conn"`
    MaxLifetime int      `yaml:"max_lifetime"` // 分钟
}

type Database struct {
    Master     *gorm.DB
    Slaves     []*gorm.DB
    current    int
    mu         sync.RWMutex
    healthChan chan bool
}

var (
    instance *Database
    once     sync.Once
)

func NewDatabase(config DBConfig) (*Database, error) {
    var err error
    once.Do(func() {
        instance = &Database{
            healthChan: make(chan bool, 1),
        }
        err = instance.init(config)
    })
    
    if err != nil {
        return nil, fmt.Errorf("failed to initialize database: %w", err)
    }
    
    // 启动健康检查
    go instance.healthCheck()
    
    return instance, nil
}

func (db *Database) init(config DBConfig) error {
    // 连接主库
    master, err := gorm.Open(mysql.Open(config.Master), &gorm.Config{
        Logger: logger.Default.LogMode(logger.Info),
        NamingStrategy: schema.NamingStrategy{
            SingularTable: true,
        },
    })
    if err != nil {
        return fmt.Errorf("failed to connect master: %w", err)
    }
    
    sqlDB, err := master.DB()
    if err == nil {
        sqlDB.SetMaxIdleConns(config.MaxIdleConn)
        sqlDB.SetMaxOpenConns(config.MaxOpenConn)
        sqlDB.SetConnMaxLifetime(time.Duration(config.MaxLifetime) * time.Minute)
    }
    
    db.Master = master
    
    // 连接从库
    for i, slaveDSN := range config.Slaves {
        slave, err := gorm.Open(mysql.Open(slaveDSN), &gorm.Config{
            Logger: logger.Default.LogMode(logger.Info),
            NamingStrategy: schema.NamingStrategy{
                SingularTable: true,
            },
        })
        if err != nil {
            return fmt.Errorf("failed to connect slave %d: %w", i, err)
        }
        
        sqlDB, err := slave.DB()
        if err == nil {
            sqlDB.SetMaxIdleConns(config.MaxIdleConn)
            sqlDB.SetMaxOpenConns(config.MaxOpenConn)
            sqlDB.SetConnMaxLifetime(time.Duration(config.MaxLifetime) * time.Minute)
        }
        
        db.Slaves = append(db.Slaves, slave)
    }
    
    return nil
}

// 获取从库(轮询负载均衡)
func (db *Database) GetSlave() *gorm.DB {
    db.mu.RLock()
    defer db.mu.RUnlock()
    
    if len(db.Slaves) == 0 {
        return db.Master
    }
    
    slave := db.Slaves[db.current]
    db.current = (db.current + 1) % len(db.Slaves)
    return slave
}

// 健康检查
func (db *Database) healthCheck() {
    ticker := time.NewTicker(30 * time.Second)
    defer ticker.Stop()
    
    for {
        select {
        case <-ticker.C:
            db.checkSlaveHealth()
        case <-db.healthChan:
            return
        }
    }
}

func (db *Database) checkSlaveHealth() {
    db.mu.Lock()
    defer db.mu.Unlock()
    
    healthySlaves := make([]*gorm.DB, 0)
    
    for i, slave := range db.Slaves {
        sqlDB, err := slave.DB()
        if err != nil {
            fmt.Printf("Slave %d: failed to get sql.DB: %v\n", i, err)
            continue
        }
        
        ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
        err = sqlDB.PingContext(ctx)
        cancel()
        
        if err != nil {
            fmt.Printf("Slave %d is unhealthy: %v\n", i, err)
        } else {
            healthySlaves = append(healthySlaves, slave)
        }
    }
    
    db.Slaves = healthySlaves
    if db.current >= len(db.Slaves) {
        db.current = 0
    }
}

func (db *Database) Close() {
    close(db.healthChan)
    
    if db.Master != nil {
        sqlDB, _ := db.Master.DB()
        sqlDB.Close()
    }
    
    for _, slave := range db.Slaves {
        sqlDB, _ := slave.DB()
        sqlDB.Close()
    }
}

2. 增强的读写分离中间件

// internal/middleware/db_router.go
package middleware

import (
    "github.com/gin-gonic/gin"
    "net/http"
    "regexp"
    "strings"
    "time"
    
    "your-project/pkg/database"
    "your-project/pkg/logger"
)

type DBType string

const (
    DBTypeWrite DBType = "write"
    DBTypeRead  DBType = "read"
    DBTypeAuto  DBType = "auto"
)

type DBContext struct {
    db *database.Database
}

func NewDBContext(db *database.Database) *DBContext {
    return &DBContext{db: db}
}

// 白名单:这些路径强制走主库
var masterOnlyPaths = []string{
    "/api/v1/users/profile",
    "/api/v1/orders/checkout",
}

// 黑名单:这些路径强制走从库
var slaveOnlyPaths = []string{
    "/api/v1/reports/",
    "/api/v1/analytics/",
}

func (ctx *DBContext) Middleware() gin.HandlerFunc {
    return func(c *gin.Context) {
        start := time.Now()
        
        // 获取数据库连接
        dbType := ctx.determineDBType(c)
        var db *gorm.DB
        
        switch dbType {
        case DBTypeWrite:
            db = ctx.db.Master
            c.Set("db_type", "master")
        case DBTypeRead:
            db = ctx.db.GetSlave()
            c.Set("db_type", "slave")
        default:
            // 默认根据HTTP方法决定
            if c.Request.Method == http.MethodGet || c.Request.Method == http.MethodHead {
                db = ctx.db.GetSlave()
                c.Set("db_type", "slave")
            } else {
                db = ctx.db.Master
                c.Set("db_type", "master")
            }
        }
        
        c.Set("db", db)
        
        c.Next()
        
        // 记录查询耗时
        latency := time.Since(start)
        logger.Infof("DB Route | Method: %s | Path: %s | DB Type: %s | Latency: %v",
            c.Request.Method,
            c.Request.URL.Path,
            c.GetString("db_type"),
            latency)
    }
}

func (ctx *DBContext) determineDBType(c *gin.Context) DBType {
    path := c.Request.URL.Path
    method := c.Request.Method
    
    // 检查白名单
    for _, p := range masterOnlyPaths {
        if strings.HasPrefix(path, p) {
            return DBTypeWrite
        }
    }
    
    // 检查黑名单
    for _, p := range slaveOnlyPaths {
        if strings.HasPrefix(path, p) {
            return DBTypeRead
        }
    }
    
    // 特殊查询参数
    if c.Query("force_master") == "true" {
        return DBTypeWrite
    }
    
    // 根据方法决定
    if method == http.MethodPost || method == http.MethodPut || 
       method == http.MethodPatch || method == http.MethodDelete {
        return DBTypeWrite
    }
    
    return DBTypeRead
}

// 获取数据库连接
func GetDB(c *gin.Context) *gorm.DB {
    if db, exists := c.Get("db"); exists {
        return db.(*gorm.DB)
    }
    return nil
}

3. Repository模式实现

// internal/repository/base_repo.go
package repository

import (
    "gorm.io/gorm"
    "github.com/gin-gonic/gin"
)

type BaseRepository struct {
    db *gorm.DB
}

func NewBaseRepository(db *gorm.DB) *BaseRepository {
    return &BaseRepository{db: db}
}

func (r *BaseRepository) DB() *gorm.DB {
    return r.db
}

func (r *BaseRepository) WithContext(c *gin.Context) *gorm.DB {
    if db := GetDB(c); db != nil {
        return db
    }
    return r.db
}
// internal/repository/user_repo.go
package repository

import (
    "context"
    "time"
    
    "gorm.io/gorm"
    "your-project/internal/model"
)

type UserRepository interface {
    Create(c *gin.Context, user *model.User) error
    Update(c *gin.Context, user *model.User) error
    Delete(c *gin.Context, id uint) error
    FindByID(c *gin.Context, id uint) (*model.User, error)
    FindByEmail(c *gin.Context, email string) (*model.User, error)
    FindAll(c *gin.Context, page, pageSize int) ([]model.User, int64, error)
    Count(c *gin.Context) (int64, error)
}

type userRepository struct {
    BaseRepository
}

func NewUserRepository(db *gorm.DB) UserRepository {
    return &userRepository{BaseRepository: *NewBaseRepository(db)}
}

func (r *userRepository) Create(c *gin.Context, user *model.User) error {
    return r.WithContext(c).Create(user).Error
}

func (r *userRepository) Update(c *gin.Context, user *model.User) error {
    return r.WithContext(c).Save(user).Error
}

func (r *userRepository) Delete(c *gin.Context, id uint) error {
    return r.WithContext(c).Delete(&model.User{}, id).Error
}

func (r *userRepository) FindByID(c *gin.Context, id uint) (*model.User, error) {
    var user model.User
    err := r.WithContext(c).First(&user, id).Error
    return &user, err
}

func (r *userRepository) FindByEmail(c *gin.Context, email string) (*model.User, error) {
    var user model.User
    err := r.WithContext(c).Where("email = ?", email).First(&user).Error
    return &user, err
}

func (r *userRepository) FindAll(c *gin.Context, page, pageSize int) ([]model.User, int64, error) {
    var users []model.User
    var total int64
    
    db := r.WithContext(c)
    
    // 获取总数
    if err := db.Model(&model.User{}).Count(&total).Error; err != nil {
        return nil, 0, err
    }
    
    // 分页查询
    offset := (page - 1) * pageSize
    err := db.Offset(offset).Limit(pageSize).Find(&users).Error
    
    return users, total, err
}

func (r *userRepository) Count(c *gin.Context) (int64, error) {
    var count int64
    err := r.WithContext(c).Model(&model.User{}).Count(&count).Error
    return count, err
}

4. Service层

// internal/service/user_service.go
package service

import (
    "errors"
    "time"
    
    "github.com/gin-gonic/gin"
    "your-project/internal/model"
    "your-project/internal/repository"
)

type UserService interface {
    CreateUser(c *gin.Context, req *CreateUserRequest) (*model.User, error)
    UpdateUser(c *gin.Context, id uint, req *UpdateUserRequest) (*model.User, error)
    GetUser(c *gin.Context, id uint) (*model.User, error)
    GetUsers(c *gin.Context, page, pageSize int) (*UserListResponse, error)
    DeleteUser(c *gin.Context, id uint) error
    ForceMasterQuery(c *gin.Context, id uint) (*model.User, error)
}

type userService struct {
    userRepo repository.UserRepository
}

func NewUserService(userRepo repository.UserRepository) UserService {
    return &userService{userRepo: userRepo}
}

type CreateUserRequest struct {
    Username string `json:"username" binding:"required,min=3,max=50"`
    Email    string `json:"email" binding:"required,email"`
}

type UpdateUserRequest struct {
    Username string `json:"username" binding:"omitempty,min=3,max=50"`
    Email    string `json:"email" binding:"omitempty,email"`
}

type UserListResponse struct {
    Users []model.User `json:"users"`
    Total int64        `json:"total"`
    Page  int          `json:"page"`
    Size  int          `json:"size"`
}

func (s *userService) CreateUser(c *gin.Context, req *CreateUserRequest) (*model.User, error) {
    // 检查邮箱是否已存在
    existing, _ := s.userRepo.FindByEmail(c, req.Email)
    if existing != nil && existing.ID > 0 {
        return nil, errors.New("email already exists")
    }
    
    user := &model.User{
        Username: req.Username,
        Email:    req.Email,
    }
    
    if err := s.userRepo.Create(c, user); err != nil {
        return nil, err
    }
    
    return user, nil
}

func (s *userService) UpdateUser(c *gin.Context, id uint, req *UpdateUserRequest) (*model.User, error) {
    user, err := s.userRepo.FindByID(c, id)
    if err != nil {
        return nil, errors.New("user not found")
    }
    
    if req.Username != "" {
        user.Username = req.Username
    }
    if req.Email != "" && req.Email != user.Email {
        // 检查新邮箱是否已被使用
        existing, _ := s.userRepo.FindByEmail(c, req.Email)
        if existing != nil && existing.ID != user.ID {
            return nil, errors.New("email already in use")
        }
        user.Email = req.Email
    }
    
    if err := s.userRepo.Update(c, user); err != nil {
        return nil, err
    }
    
    return user, nil
}

func (s *userService) GetUser(c *gin.Context, id uint) (*model.User, error) {
    return s.userRepo.FindByID(c, id)
}

func (s *userService) GetUsers(c *gin.Context, page, pageSize int) (*UserListResponse, error) {
    if page < 1 {
        page = 1
    }
    if pageSize < 1 {
        pageSize = 10
    }
    if pageSize > 100 {
        pageSize = 100
    }
    
    users, total, err := s.userRepo.FindAll(c, page, pageSize)
    if err != nil {
        return nil, err
    }
    
    return &UserListResponse{
        Users: users,
        Total: total,
        Page:  page,
        Size:  len(users),
    }, nil
}

func (s *userService) DeleteUser(c *gin.Context, id uint) error {
    return s.userRepo.Delete(c, id)
}

// 强制从主库读取(用于需要强一致性的场景)
func (s *userService) ForceMasterQuery(c *gin.Context, id uint) (*model.User, error) {
    // 复制上下文,添加强制主库查询参数
    newCtx := &gin.Context{}
    *newCtx = *c
    newCtx.Request.URL.RawQuery = "force_master=true"
    
    return s.userRepo.FindByID(newCtx, id)
}

5. Handler层

// internal/handler/user_handler.go
package handler

import (
    "net/http"
    "strconv"
    
    "github.com/gin-gonic/gin"
    "your-project/internal/service"
)

type UserHandler struct {
    userService service.UserService
}

func NewUserHandler(userService service.UserService) *UserHandler {
    return &UserHandler{userService: userService}
}

func (h *UserHandler) RegisterRoutes(router *gin.RouterGroup) {
    users := router.Group("/users")
    {
        users.POST("", h.CreateUser)
        users.GET("", h.GetUsers)
        users.GET("/:id", h.GetUser)
        users.PUT("/:id", h.UpdateUser)
        users.DELETE("/:id", h.DeleteUser)
        users.GET("/:id/force-master", h.ForceMasterQuery)
    }
}

func (h *UserHandler) CreateUser(c *gin.Context) {
    var req service.CreateUserRequest
    if err := c.ShouldBindJSON(&req); err != nil {
        c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()})
        return
    }
    
    user, err := h.userService.CreateUser(c, &req)
    if err != nil {
        c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()})
        return
    }
    
    c.JSON(http.StatusCreated, user)
}

func (h *UserHandler) GetUsers(c *gin.Context) {
    page, _ := strconv.Atoi(c.DefaultQuery("page", "1"))
    pageSize, _ := strconv.Atoi(c.DefaultQuery("page_size", "10"))
    
    resp, err := h.userService.GetUsers(c, page, pageSize)
    if err != nil {
        c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
        return
    }
    
    c.JSON(http.StatusOK, resp)
}

func (h *UserHandler) GetUser(c *gin.Context) {
    id, err := strconv.ParseUint(c.Param("id"), 10, 32)
    if err != nil {
        c.JSON(http.StatusBadRequest, gin.H{"error": "invalid id"})
        return
    }
    
    user, err := h.userService.GetUser(c, uint(id))
    if err != nil {
        c.JSON(http.StatusNotFound, gin.H{"error": "user not found"})
        return
    }
    
    c.JSON(http.StatusOK, user)
}

func (h *UserHandler) UpdateUser(c *gin.Context) {
    id, err := strconv.ParseUint(c.Param("id"), 10, 32)
    if err != nil {
        c.JSON(http.StatusBadRequest, gin.H{"error": "invalid id"})
        return
    }
    
    var req service.UpdateUserRequest
    if err := c.ShouldBindJSON(&req); err != nil {
        c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()})
        return
    }
    
    user, err := h.userService.UpdateUser(c, uint(id), &req)
    if err != nil {
        c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()})
        return
    }
    
    c.JSON(http.StatusOK, user)
}

func (h *UserHandler) DeleteUser(c *gin.Context) {
    id, err := strconv.ParseUint(c.Param("id"), 10, 32)
    if err != nil {
        c.JSON(http.StatusBadRequest, gin.H{"error": "invalid id"})
        return
    }
    
    if err := h.userService.DeleteUser(c, uint(id)); err != nil {
        c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
        return
    }
    
    c.JSON(http.StatusNoContent, nil)
}

func (h *UserHandler) ForceMasterQuery(c *gin.Context) {
    id, err := strconv.ParseUint(c.Param("id"), 10, 32)
    if err != nil {
        c.JSON(http.StatusBadRequest, gin.H{"error": "invalid id"})
        return
    }
    
    user, err := h.userService.ForceMasterQuery(c, uint(id))
    if err != nil {
        c.JSON(http.StatusNotFound, gin.H{"error": "user not found"})
        return
    }
    
    c.JSON(http.StatusOK, gin.H{
        "user": user,
        "note": "This query was forced to use master database",
    })
}

6. 主程序入口

// cmd/server/main.go
package main

import (
    "context"
    "fmt"
    "log"
    "net/http"
    "os"
    "os/signal"
    "syscall"
    "time"
    
    "github.com/gin-gonic/gin"
    "gopkg.in/yaml.v3"
    
    "your-project/config"
    "your-project/internal/handler"
    "your-project/internal/middleware"
    "your-project/internal/repository"
    "your-project/internal/router"
    "your-project/internal/service"
    "your-project/pkg/database"
    "your-project/pkg/logger"
)

func main() {
    // 加载配置
    cfg, err := loadConfig("config/config.yaml")
    if err != nil {
        log.Fatal("Failed to load config:", err)
    }
    
    // 初始化日志
    logger.InitLogger(cfg.Server.Mode)
    
    // 初始化数据库
    db, err := database.NewDatabase(cfg.Database)
    if err != nil {
        log.Fatal("Failed to connect database:", err)
    }
    defer db.Close()
    
    // 初始化仓库
    userRepo := repository.NewUserRepository(db.Master)
    productRepo := repository.NewProductRepository(db.Master)
    
    // 初始化服务
    userService := service.NewUserService(userRepo)
    productService := service.NewProductService(productRepo)
    healthService := service.NewHealthService(db)
    
    // 初始化处理器
    userHandler := handler.NewUserHandler(userService)
    productHandler := handler.NewProductHandler(productService)
    healthHandler := handler.NewHealthHandler(healthService)
    
    // 创建Gin应用
    app := gin.New()
    
    // 中间件
    app.Use(gin.Recovery())
    app.Use(middleware.CORSMiddleware())
    app.Use(middleware.LoggerMiddleware())
    app.Use(middleware.NewDBContext(db).Middleware())
    
    // 注册路由
    apiRouter := router.NewRouter(userHandler, productHandler, healthHandler)
    apiRouter.Register(app)
    
    // 启动服务器
    srv := &http.Server{
        Addr:    fmt.Sprintf(":%d", cfg.Server.Port),
        Handler: app,
    }
    
    go func() {
        logger.Infof("Server starting on port %d", cfg.Server.Port)
        if err := srv.ListenAndServe(); err != nil && err != http.ErrServerClosed {
            logger.Fatalf("Server failed: %v", err)
        }
    }()
    
    // 优雅关闭
    quit := make(chan os.Signal, 1)
    signal.Notify(quit, syscall.SIGINT, syscall.SIGTERM)
    <-quit
    
    logger.Info("Shutting down server...")
    
    ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
    defer cancel()
    
    if err := srv.Shutdown(ctx); err != nil {
        logger.Fatalf("Server forced to shutdown: %v", err)
    }
    
    logger.Info("Server exiting")
}

func loadConfig(path string) (*config.Config, error) {
    data, err := os.ReadFile(path)
    if err != nil {
        return nil, err
    }
    
    var cfg config.Config
    if err := yaml.Unmarshal(data, &cfg); err != nil {
        return nil, err
    }
    
    return &cfg, nil
}

7. 配置文件



# config/config.yaml
server:
  port: 8080
  mode: "debug"
  read_timeout: 10
  write_timeout: 10

database:
  master: "app_user:app_pass123@tcp(localhost:3307)/app_db?charset=utf8mb4&parseTime=True&loc=Local"
  slaves:
    - "root:root123@tcp(localhost:3308)/app_db?charset=utf8mb4&parseTime=True&loc=Local"
    - "root:root123@tcp(localhost:3309)/app_db?charset=utf8mb4&parseTime=True&loc=Local"
  max_idle_conn: 10
  max_open_conn: 100
  max_lifetime: 60

redis:
  addr: "localhost:6379"
  password: ""
  db: 0

8. 监控和健康检查端点

// internal/handler/health_handler.go
package handler

import (
    "net/http"
    
    "github.com/gin-gonic/gin"
    "your-project/internal/service"
)

type HealthHandler struct {
    healthService service.HealthService
}

func NewHealthHandler(healthService service.HealthService) *HealthHandler {
    return &HealthHandler{healthService: healthService}
}

func (h *HealthHandler) RegisterRoutes(router *gin.RouterGroup) {
    router.GET("/health", h.Health)
    router.GET("/health/db", h.DatabaseHealth)
    router.GET("/health/replication", h.ReplicationStatus)
}

func (h *HealthHandler) Health(c *gin.Context) {
    c.JSON(http.StatusOK, gin.H{
        "status":  "healthy",
        "service": "gin-mysql-replication",
    })
}

func (h *HealthHandler) DatabaseHealth(c *gin.Context) {
    masterStatus, slaveStatus, err := h.healthService.CheckDatabaseHealth(c)
    if err != nil {
        c.JSON(http.StatusServiceUnavailable, gin.H{
            "status":  "unhealthy",
            "error":   err.Error(),
            "master":  masterStatus,
            "slaves":  slaveStatus,
        })
        return
    }
    
    c.JSON(http.StatusOK, gin.H{
        "status":  "healthy",
        "master":  masterStatus,
        "slaves":  slaveStatus,
    })
}

func (h *HealthHandler) ReplicationStatus(c *gin.Context) {
    status, err := h.healthService.GetReplicationStatus(c)
    if err != nil {
        c.JSON(http.StatusInternalServerError, gin.H{
            "status":  "error",
            "error":   err.Error(),
        })
        return
    }
    
    c.JSON(http.StatusOK, gin.H{
        "status":  "success",
        "data":    status,
    })
}

三、测试脚本

#!/bin/bash
# test-replication.sh

echo "=== 测试MySQL主从复制 ==="

# 1. 写入主库
echo "1. 写入主库..."
curl -X POST http://localhost:8080/api/v1/users \
  -H "Content-Type: application/json" \
  -d '{"username": "testuser1", "email": "test1@example.com"}'

curl -X POST http://localhost:8080/api/v1/users \
  -H "Content-Type: application/json" \
  -d '{"username": "testuser2", "email": "test2@example.com"}'

# 2. 从从库读取
echo -e "\n2. 从从库读取..."
curl -X GET http://localhost:8080/api/v1/users

# 3. 查看复制状态
echo -e "\n3. 查看复制状态..."
curl -X GET http://localhost:8080/api/v1/health/replication

# 4. 强制从主库读取
echo -e "\n4. 强制从主库读取..."
curl -X GET http://localhost:8080/api/v1/users/1/force-master

# 5. 压力测试(可选)
echo -e "\n5. 运行压力测试..."
echo "使用wrk进行压力测试:"
echo "wrk -t12 -c400 -d30s http://localhost:8080/api/v1/users"

四、性能监控



// internal/monitoring/metrics.go
package monitoring

import (
    "time"
    
    "github.com/prometheus/client_golang/prometheus"
    "github.com/prometheus/client_golang/prometheus/promauto"
)

var (
    dbQueryDuration = promauto.NewHistogramVec(prometheus.HistogramOpts{
        Name:    "db_query_duration_seconds",
        Help:    "Database query duration in seconds",
        Buckets: []float64{.005, .01, .025, .05, .1, .25, .5, 1},
    }, []string{"type", "operation"})
    
    dbQueryTotal = promauto.NewCounterVec(prometheus.CounterOpts{
        Name: "db_query_total",
        Help: "Total number of database queries",
    }, []string{"type", "operation"})
    
    dbErrorsTotal = promauto.NewCounterVec(prometheus.CounterOpts{
        Name: "db_errors_total",
        Help: "Total number of database errors",
    }, []string{"type"})
)

func RecordDBQuery(dbType, operation string, duration time.Duration) {
    dbQueryDuration.WithLabelValues(dbType, operation).Observe(duration.Seconds())
    dbQueryTotal.WithLabelValues(dbType, operation).Inc()
}

func RecordDBError(dbType string) {
    dbErrorsTotal.WithLabelValues(dbType).Inc()
}

五、部署说明

1. 构建和运行

# 1. 克隆项目
git clone <your-repo>
cd gin-mysql-replication

# 2. 启动数据库
docker-compose up -d

# 3. 设置主从复制
chmod +x setup-replication.sh
./setup-replication.sh

# 4. 启动应用
go mod tidy
go run cmd/server/main.go

# 5. 或使用Docker部署应用
docker build -t gin-app .
docker run -p 8080:8080 --network mysql-replication-net gin-app

2. Dockerfile

# Dockerfile
FROM golang:1.21-alpine AS builder

WORKDIR /app
COPY go.mod go.sum ./
RUN go mod download
COPY . .
RUN CGO_ENABLED=0 GOOS=linux go build -o main cmd/server/main.go

FROM alpine:latest
RUN apk --no-cache add ca-certificates
WORKDIR /root/
COPY --from=builder /app/main .
COPY --from=builder /app/config/config.yaml ./config/
EXPOSE 8080
CMD ["./main"]

六、总结

这个完整示例展示了:

  1. Docker搭建MySQL主从集群:一键部署主库+2个从库

  2. 自动配置复制:通过脚本自动设置主从关系

  3. 智能读写分离:基于HTTP方法和路径的自动路由

  4. 健康检查:数据库连接监控和故障处理

  5. 优雅架构:Repository-Service-Handler分层

  6. 监控指标:Prometheus指标收集

  7. 生产就绪:连接池、超时、优雅关闭等

关键特性:

  • ✅ 自动读写分离

  • ✅ 连接池管理

  • ✅ 健康检查和故障恢复

  • ✅ 监控和指标

  • ✅ 可配置的路由策略

  • ✅ 支持强制主库查询

  • ✅ 优雅的API设计

这个架构可以轻松扩展到多个从库,支持负载均衡和高可用性。

gin框架 gorm 的Mysql 实现基于binlog的主从同步

在Go语言中使用Gin、GORM和MySQL的binlog实现主从同步,需要结合多个组件。以下是完整的实现方案:

1. 项目结构

├── config/
│   ├── config.go
│   └── database.go
├── models/
│   └── model.go
├── replication/
│   ├── binlog_reader.go
│   ├── replication_manager.go
│   └── event_handler.go
├── services/
│   └── sync_service.go
├── routers/
│   └── router.go
└── main.go

2. 数据库配置

// config/database.go
package config

import (
	"gorm.io/driver/mysql"
	"gorm.io/gorm"
	"gorm.io/gorm/logger"
	"log"
	"time"
)

type DatabaseConfig struct {
	MasterDSN string
	SlaveDSN  string
}

var (
	MasterDB *gorm.DB
	SlaveDB  *gorm.DB
)

func InitDatabases(cfg DatabaseConfig) error {
	var err error
	
	// 主数据库
	MasterDB, err = gorm.Open(mysql.Open(cfg.MasterDSN), &gorm.Config{
		Logger: logger.Default.LogMode(logger.Info),
	})
	if err != nil {
		return err
	}
	
	// 从数据库
	SlaveDB, err = gorm.Open(mysql.Open(cfg.SlaveDSN), &gorm.Config{
		Logger: logger.Default.LogMode(logger.Info),
	})
	if err != nil {
		return err
	}
	
	// 配置连接池
	sqlDB, err := MasterDB.DB()
	if err != nil {
		return err
	}
	sqlDB.SetMaxIdleConns(10)
	sqlDB.SetMaxOpenConns(100)
	sqlDB.SetConnMaxLifetime(time.Hour)
	
	slaveSqlDB, err := SlaveDB.DB()
	if err != nil {
		return err
	}
	slaveSqlDB.SetMaxIdleConns(10)
	slaveSqlDB.SetMaxOpenConns(100)
	slaveSqlDB.SetConnMaxLifetime(time.Hour)
	
	return nil
}

3. 数据模型

// models/model.go
package models

import (
	"gorm.io/gorm"
	"time"
)

type User struct {
	ID        uint           `gorm:"primarykey" json:"id"`
	Name      string         `gorm:"size:100" json:"name"`
	Email     string         `gorm:"size:100;uniqueIndex" json:"email"`
	Age       int            `json:"age"`
	CreatedAt time.Time      `json:"created_at"`
	UpdatedAt time.Time      `json:"updated_at"`
	DeletedAt gorm.DeletedAt `gorm:"index" json:"-"`
}

type Product struct {
	ID        uint           `gorm:"primarykey" json:"id"`
	Name      string         `gorm:"size:200" json:"name"`
	Price     float64        `json:"price"`
	Stock     int            `json:"stock"`
	CreatedAt time.Time      `json:"created_at"`
	UpdatedAt time.Time      `json:"updated_at"`
	DeletedAt gorm.DeletedAt `gorm:"index" json:"-"`
}

// BinlogPosition 记录binlog位置
type BinlogPosition struct {
	ID            uint   `gorm:"primarykey"`
	BinlogName    string `gorm:"size:120;index"`
	BinlogPos     uint32
	GtidSet       string `gorm:"type:text"`
	LastSyncedAt  time.Time
	CreatedAt     time.Time
	UpdatedAt     time.Time
}

4. Binlog读取器(使用go-mysql)

首先安装依赖:

go get github.com/go-mysql-org/go-mysql
go get github.com/siddontang/go-log/log
// replication/binlog_reader.go
package replication

import (
	"context"
	"encoding/json"
	"fmt"
	"log"
	"time"
	
	"github.com/go-mysql-org/go-mysql/canal"
	"github.com/go-mysql-org/go-mysql/mysql"
	"github.com/go-mysql-org/go-mysql/replication"
	"gorm.io/gorm"
)

type BinlogEventHandler interface {
	OnRow(e *canal.RowsEvent) error
	String() string
}

type BinlogReader struct {
	canal    *canal.Canal
	handlers map[string]BinlogEventHandler
	position *mysql.Position
	gtidSet  mysql.GTIDSet
	db       *gorm.DB
	cfg      *canal.Config
}

type BinlogConfig struct {
	Host     string
	Port     int
	User     string
	Password string
	ServerID uint32
}

func NewBinlogReader(cfg BinlogConfig, db *gorm.DB) (*BinlogReader, error) {
	canalCfg := canal.NewDefaultConfig()
	canalCfg.Addr = fmt.Sprintf("%s:%d", cfg.Host, cfg.Port)
	canalCfg.User = cfg.User
	canalCfg.Password = cfg.Password
	canalCfg.ServerID = cfg.ServerID
	canalCfg.Dump.ExecutionPath = "" // 不使用mysqldump
	
	// 只监听特定的数据库和表
	canalCfg.IncludeTableRegex = []string{"test\\..*"}
	
	c, err := canal.NewCanal(canalCfg)
	if err != nil {
		return nil, err
	}
	
	return &BinlogReader{
		canal:    c,
		handlers: make(map[string]BinlogEventHandler),
		db:       db,
		cfg:      canalCfg,
	}, nil
}

func (r *BinlogReader) RegisterHandler(table string, handler BinlogEventHandler) {
	r.handlers[table] = handler
}

func (r *BinlogReader) LoadPosition() error {
	var pos BinlogPosition
	if err := r.db.Last(&pos).Error; err != nil {
		if err == gorm.ErrRecordNotFound {
			// 第一次启动,从当前位点开始
			pos, err := r.canal.GetMasterPos()
			if err != nil {
				return err
			}
			r.position = &pos
			return r.savePosition(pos.Name, pos.Pos, "")
		}
		return err
	}
	
	r.position = &mysql.Position{
		Name: pos.BinlogName,
		Pos:  pos.BinlogPos,
	}
	
	if pos.GtidSet != "" {
		gtidSet, err := mysql.ParseGTIDSet("mysql", pos.GtidSet)
		if err != nil {
			return err
		}
		r.gtidSet = gtidSet
	}
	
	return nil
}

func (r *BinlogReader) savePosition(name string, pos uint32, gtidSet string) error {
	binlogPos := BinlogPosition{
		BinlogName:   name,
		BinlogPos:    pos,
		GtidSet:      gtidSet,
		LastSyncedAt: time.Now(),
	}
	
	return r.db.Save(&binlogPos).Error
}

func (r *BinlogReader) Start(ctx context.Context) error {
	// 设置事件处理器
	r.canal.SetEventHandler(&binlogEventHandler{
		reader:   r,
		handlers: r.handlers,
		db:       r.db,
	})
	
	// 加载之前的同步位置
	if err := r.LoadPosition(); err != nil {
		return err
	}
	
	// 开始同步
	if r.gtidSet != nil {
		return r.canal.StartFromGTID(r.gtidSet)
	} else if r.position != nil {
		return r.canal.RunFrom(*r.position)
	} else {
		return r.canal.Run()
	}
}

func (r *BinlogReader) Stop() {
	r.canal.Close()
}

// 自定义事件处理器
type binlogEventHandler struct {
	canal.DummyEventHandler
	reader   *BinlogReader
	handlers map[string]BinlogEventHandler
	db       *gorm.DB
}

func (h *binlogEventHandler) OnRow(e *canal.RowsEvent) error {
	log.Printf("OnRow: %s %v", e.Table.Name, e.Action)
	
	// 查找对应的处理器
	handler, ok := h.handlers[e.Table.Name]
	if ok {
		return handler.OnRow(e)
	}
	
	// 默认处理器
	return h.defaultHandler(e)
}

func (h *binlogEventHandler) defaultHandler(e *canal.RowsEvent) error {
	// 这里可以将事件存储到消息队列或直接同步到从库
	data, _ := json.Marshal(map[string]interface{}{
		"table":  e.Table.Name,
		"action": e.Action,
		"rows":   e.Rows,
		"time":   time.Now(),
	})
	
	log.Printf("Binlog Event: %s", string(data))
	return nil
}

func (h *binlogEventHandler) String() string {
	return "BinlogEventHandler"
}

func (h *binlogEventHandler) OnPosSynced(pos mysql.Position, gtid mysql.GTIDSet, force bool) error {
	// 保存同步位置
	var gtidStr string
	if gtid != nil {
		gtidStr = gtid.String()
	}
	return h.reader.savePosition(pos.Name, pos.Pos, gtidStr)
}

5. 同步管理器

// replication/replication_manager.go
package replication

import (
	"context"
	"fmt"
	"log"
	"sync"
	"time"
	
	"gorm.io/gorm"
)

type ReplicationManager struct {
	reader      *BinlogReader
	masterDB    *gorm.DB
	slaveDB     *gorm.DB
	isRunning   bool
	mu          sync.RWMutex
	ctx         context.Context
	cancel      context.CancelFunc
	errorChan   chan error
}

func NewReplicationManager(masterDB, slaveDB *gorm.DB, cfg BinlogConfig) (*ReplicationManager, error) {
	reader, err := NewBinlogReader(cfg, masterDB)
	if err != nil {
		return nil, err
	}
	
	ctx, cancel := context.WithCancel(context.Background())
	
	return &ReplicationManager{
		reader:    reader,
		masterDB:  masterDB,
		slaveDB:   slaveDB,
		ctx:       ctx,
		cancel:    cancel,
		errorChan: make(chan error, 10),
	}, nil
}

func (m *ReplicationManager) Start() error {
	m.mu.Lock()
	defer m.mu.Unlock()
	
	if m.isRunning {
		return fmt.Errorf("replication manager is already running")
	}
	
	// 注册表处理器
	m.registerHandlers()
	
	// 启动binlog读取
	go func() {
		if err := m.reader.Start(m.ctx); err != nil {
			m.errorChan <- err
		}
	}()
	
	m.isRunning = true
	log.Println("Replication manager started")
	
	return nil
}

func (m *ReplicationManager) Stop() {
	m.mu.Lock()
	defer m.mu.Unlock()
	
	if !m.isRunning {
		return
	}
	
	m.cancel()
	m.reader.Stop()
	m.isRunning = false
	
	log.Println("Replication manager stopped")
}

func (m *ReplicationManager) registerHandlers() {
	// 注册用户表处理器
	userHandler := &UserEventHandler{slaveDB: m.slaveDB}
	m.reader.RegisterHandler("users", userHandler)
	
	// 注册产品表处理器
	productHandler := &ProductEventHandler{slaveDB: m.slaveDB}
	m.reader.RegisterHandler("products", productHandler)
}

func (m *ReplicationManager) GetErrors() <-chan error {
	return m.errorChan
}

func (m *ReplicationManager) IsRunning() bool {
	m.mu.RLock()
	defer m.mu.RUnlock()
	return m.isRunning
}

// 健康检查
func (m *ReplicationManager) HealthCheck() error {
	if !m.IsRunning() {
		return fmt.Errorf("replication is not running")
	}
	
	// 检查主从延迟
	var masterPos, slavePos string
	m.masterDB.Raw("SHOW MASTER STATUS").Scan(&masterPos)
	m.slaveDB.Raw("SHOW SLAVE STATUS").Scan(&slavePos)
	
	log.Printf("Master position: %s, Slave position: %s", masterPos, slavePos)
	
	return nil
}

6. 具体表的处理器

// replication/event_handler.go
package replication

import (
	"fmt"
	"log"
	"time"
	
	"github.com/go-mysql-org/go-mysql/canal"
	"gorm.io/gorm"
	"your-project/models"
)

type UserEventHandler struct {
	slaveDB *gorm.DB
}

func (h *UserEventHandler) OnRow(e *canal.RowsEvent) error {
	switch e.Action {
	case canal.InsertAction:
		return h.handleInsert(e)
	case canal.UpdateAction:
		return h.handleUpdate(e)
	case canal.DeleteAction:
		return h.handleDelete(e)
	default:
		return fmt.Errorf("unsupported action: %s", e.Action)
	}
}

func (h *UserEventHandler) handleInsert(e *canal.RowsEvent) error {
	for _, row := range e.Rows {
		user := models.User{
			ID:        row[0].(uint32),
			Name:      row[1].(string),
			Email:     row[2].(string),
			Age:       int(row[3].(int32)),
			CreatedAt: row[4].(time.Time),
			UpdatedAt: row[5].(time.Time),
		}
		
		if err := h.slaveDB.Create(&user).Error; err != nil {
			log.Printf("Failed to insert user: %v", err)
			return err
		}
	}
	log.Printf("Inserted %d users", len(e.Rows))
	return nil
}

func (h *UserEventHandler) handleUpdate(e *canal.RowsEvent) error {
	for i := 0; i < len(e.Rows); i += 2 {
		oldRow := e.Rows[i]
		newRow := e.Rows[i+1]
		
		user := models.User{
			ID:        newRow[0].(uint32),
			Name:      newRow[1].(string),
			Email:     newRow[2].(string),
			Age:       int(newRow[3].(int32)),
			UpdatedAt: newRow[5].(time.Time),
		}
		
		if err := h.slaveDB.Save(&user).Error; err != nil {
			log.Printf("Failed to update user: %v", err)
			return err
		}
	}
	log.Printf("Updated %d users", len(e.Rows)/2)
	return nil
}

func (h *UserEventHandler) handleDelete(e *canal.RowsEvent) error {
	for _, row := range e.Rows {
		userID := row[0].(uint32)
		if err := h.slaveDB.Delete(&models.User{}, userID).Error; err != nil {
			log.Printf("Failed to delete user: %v", err)
			return err
		}
	}
	log.Printf("Deleted %d users", len(e.Rows))
	return nil
}

func (h *UserEventHandler) String() string {
	return "UserEventHandler"
}

type ProductEventHandler struct {
	slaveDB *gorm.DB
}

func (h *ProductEventHandler) OnRow(e *canal.RowsEvent) error {
	// 实现类似UserEventHandler的逻辑
	// ...
	return nil
}

func (h *ProductEventHandler) String() string {
	return "ProductEventHandler"
}

7. Gin路由和控制器

// routers/router.go
package routers

import (
	"github.com/gin-gonic/gin"
	"gorm.io/gorm"
	"your-project/services"
)

func SetupRouter(masterDB, slaveDB *gorm.DB, syncService *services.SyncService) *gin.Engine {
	r := gin.Default()
	
	// 健康检查
	r.GET("/health", func(c *gin.Context) {
		c.JSON(200, gin.H{
			"status": "ok",
		})
	})
	
	// 主库写操作
	r.POST("/users", func(c *gin.Context) {
		var user models.User
		if err := c.ShouldBindJSON(&user); err != nil {
			c.JSON(400, gin.H{"error": err.Error()})
			return
		}
		
		if err := masterDB.Create(&user).Error; err != nil {
			c.JSON(500, gin.H{"error": err.Error()})
			return
		}
		
		c.JSON(200, gin.H{
			"message": "User created",
			"data":    user,
		})
	})
	
	// 从库读操作
	r.GET("/users", func(c *gin.Context) {
		var users []models.User
		if err := slaveDB.Find(&users).Error; err != nil {
			c.JSON(500, gin.H{"error": err.Error()})
			return
		}
		c.JSON(200, users)
	})
	
	// 同步状态
	r.GET("/sync/status", func(c *gin.Context) {
		status := syncService.GetStatus()
		c.JSON(200, status)
	})
	
	// 手动触发同步
	r.POST("/sync/trigger", func(c *gin.Context) {
		if err := syncService.TriggerSync(); err != nil {
			c.JSON(500, gin.H{"error": err.Error()})
			return
		}
		c.JSON(200, gin.H{"message": "Sync triggered"})
	})
	
	return r
}

8. 主程序

// main.go
package main

import (
	"context"
	"log"
	"os"
	"os/signal"
	"syscall"
	"time"
	
	"your-project/config"
	"your-project/models"
	"your-project/replication"
	"your-project/routers"
	"your-project/services"
)

func main() {
	// 初始化配置
	dbConfig := config.DatabaseConfig{
		MasterDSN: "root:password@tcp(localhost:3306)/test?charset=utf8mb4&parseTime=True&loc=Local",
		SlaveDSN:  "root:password@tcp(localhost:3307)/test?charset=utf8mb4&parseTime=True&loc=Local",
	}
	
	// 初始化数据库
	if err := config.InitDatabases(dbConfig); err != nil {
		log.Fatal("Failed to init databases:", err)
	}
	
	// 自动迁移
	config.MasterDB.AutoMigrate(&models.User{}, &models.Product{}, &models.BinlogPosition{})
	config.SlaveDB.AutoMigrate(&models.User{}, &models.Product{})
	
	// 初始化binlog配置
	binlogCfg := replication.BinlogConfig{
		Host:     "localhost",
		Port:     3306,
		User:     "root",
		Password: "password",
		ServerID: 1001,
	}
	
	// 创建同步管理器
	repManager, err := replication.NewReplicationManager(
		config.MasterDB,
		config.SlaveDB,
		binlogCfg,
	)
	if err != nil {
		log.Fatal("Failed to create replication manager:", err)
	}
	
	// 创建同步服务
	syncService := services.NewSyncService(repManager)
	
	// 启动同步
	if err := syncService.Start(); err != nil {
		log.Fatal("Failed to start sync service:", err)
	}
	
	// 启动Gin服务器
	r := routers.SetupRouter(config.MasterDB, config.SlaveDB, syncService)
	
	// 优雅关闭
	ctx, stop := signal.NotifyContext(context.Background(), 
		syscall.SIGINT, syscall.SIGTERM)
	defer stop()
	
	go func() {
		if err := r.Run(":8080"); err != nil {
			log.Fatal("Failed to start server:", err)
		}
	}()
	
	// 监听错误
	go func() {
		for err := range repManager.GetErrors() {
			log.Printf("Replication error: %v", err)
		}
	}()
	
	// 定期健康检查
	ticker := time.NewTicker(30 * time.Second)
	defer ticker.Stop()
	
	go func() {
		for {
			select {
			case <-ticker.C:
				if err := repManager.HealthCheck(); err != nil {
					log.Printf("Health check failed: %v", err)
				}
			case <-ctx.Done():
				return
			}
		}
	}()
	
	// 等待终止信号
	<-ctx.Done()
	log.Println("Shutting down...")
	
	// 停止同步
	syncService.Stop()
	
	// 等待清理
	time.Sleep(2 * time.Second)
	log.Println("Server stopped")
}

9. 配置MySQL binlog

在MySQL配置文件(my.cnf)中添加:

[mysqld]
# 启用binlog
log-bin=mysql-bin
# 设置binlog格式为ROW模式
binlog-format=ROW
# 设置server-id
server-id=1
# 设置要复制的数据库
binlog-do-db=test
# 设置最大binlog大小
max_binlog_size=100M
# 设置binlog过期时间
expire_logs_days=7

10. Docker Compose 配置

version: '3.8'
services:
  mysql-master:
    image: mysql:8.0
    environment:
      MYSQL_ROOT_PASSWORD: password
      MYSQL_DATABASE: test
    ports:
      - "3306:3306"
    volumes:
      - ./master.cnf:/etc/mysql/conf.d/master.cnf
      - mysql-master-data:/var/lib/mysql
  
  mysql-slave:
    image: mysql:8.0
    environment:
      MYSQL_ROOT_PASSWORD: password
      MYSQL_DATABASE: test
    ports:
      - "3307:3306"
    volumes:
      - ./slave.cnf:/etc/mysql/conf.d/slave.cnf
      - mysql-slave-data:/var/lib/mysql
  
  app:
    build: .
    ports:
      - "8080:8080"
    depends_on:
      - mysql-master
      - mysql-slave
    environment:
      MASTER_DSN: "root:password@tcp(mysql-master:3306)/test?charset=utf8mb4&parseTime=True&loc=Local"
      SLAVE_DSN: "root:password@tcp(mysql-slave:3307)/test?charset=utf8mb4&parseTime=True&loc=Local"

volumes:
  mysql-master-data:
  mysql-slave-data:

注意事项

  • 性能考虑

    • 批量处理binlog事件

    • 使用连接池

    • 异步处理耗时操作

  • 错误处理

    • 实现重试机制

    • 记录失败事件

    • 提供手动修复功能

  • 监控

    • 监控同步延迟

    • 监控错误率

    • 监控内存使用

  • 高可用

    • 支持多从节点

    • 支持故障转移

    • 支持断点续传

这个实现提供了完整的基于binlog的主从同步方案,支持实时数据同步、断点续传、错误处理等功能。

Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐