Mysql 几种主从复制方式以及使用docker搭建主从复制架构&gin框架项目实践
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"]
六、总结
这个完整示例展示了:
-
Docker搭建MySQL主从集群:一键部署主库+2个从库
-
自动配置复制:通过脚本自动设置主从关系
-
智能读写分离:基于HTTP方法和路径的自动路由
-
健康检查:数据库连接监控和故障处理
-
优雅架构:Repository-Service-Handler分层
-
监控指标:Prometheus指标收集
-
生产就绪:连接池、超时、优雅关闭等
关键特性:
-
✅ 自动读写分离
-
✅ 连接池管理
-
✅ 健康检查和故障恢复
-
✅ 监控和指标
-
✅ 可配置的路由策略
-
✅ 支持强制主库查询
-
✅ 优雅的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的主从同步方案,支持实时数据同步、断点续传、错误处理等功能。
更多推荐



所有评论(0)