返回首页

Golang项目实践:4.3 数据库与存储

数据库与存储层决定了一个 Go 后端项目的性能上限、可维护性以及扩展能力。真正的线上系统通常不会只依赖一种存储:关系型数据库负责核心事务数据,Redis 承担热点缓存,MongoDB 存放非结构化文档,Elasticsearch 负责全文检索与复杂搜索。

这一章我们用一个完整可运行的 Go 示例项目,把下面几个主题串起来:

  • 关系型数据库:GORM / sqlx 实战
  • 数据库连接池配置与优化
  • 数据库迁移(Migration)
  • 缓存:Redis 客户端使用与缓存策略
  • NoSQL:MongoDB、Elasticsearch 集成
  • 数据访问层设计与事务管理

1. 技术选型思路

在实际工程里,不同组件的职责应该尽量清晰。

组件 适用场景 本文中的职责
MySQL 核心业务数据、事务一致性 用户与订单主数据
GORM 常规 CRUD、模型映射、事务封装 用户与订单写入
sqlx 复杂查询、报表 SQL、精细控制 用户消费统计查询
Redis 热点缓存、减轻数据库压力 用户详情缓存
MongoDB 文档型数据、结构灵活 用户画像文档
Elasticsearch 搜索、全文检索、筛选排序 用户搜索索引

一个很实用的原则是:

  • 写操作与核心事务,优先放在关系型数据库中完成。
  • 复杂统计查询,适合用 sqlx 手写 SQL。
  • 缓存、搜索、画像,采用外围存储做能力增强。
  • 跨存储一致性,优先采用“主库事务成功后,再异步或补偿同步”的思路,而不是盲目追求分布式强事务。

2. 示例项目说明

本文给出的示例项目结构如下:

  • docker-compose.yml:启动 MySQL / Redis / MongoDB / Elasticsearch
  • go.mod:依赖定义
  • migrations/:数据库迁移脚本
  • main.go:完整可运行的示例程序

这个示例程序会完成以下事情:

  1. 连接 MySQL、Redis、MongoDB、Elasticsearch
  2. 使用 Migration 初始化表结构
  3. 使用 GORM 在事务中创建用户和首单
  4. 使用 Redis 做用户详情缓存
  5. 使用 sqlx 查询用户消费统计
  6. 将用户画像写入 MongoDB
  7. 将用户索引写入 Elasticsearch 并执行搜索

3. 先启动依赖服务

3.1 docker-compose.yml

version: "3.9"

services:
  mysql:
    image: mysql:8.4
    container_name: go_storage_mysql
    restart: unless-stopped
    environment:
      MYSQL_ROOT_PASSWORD: root
      MYSQL_DATABASE: appdb
    ports:
      - "3306:3306"
    command:
      - --default-authentication-plugin=mysql_native_password
      - --character-set-server=utf8mb4
      - --collation-server=utf8mb4_unicode_ci

  redis:
    image: redis:7.2
    container_name: go_storage_redis
    restart: unless-stopped
    ports:
      - "6379:6379"

  mongo:
    image: mongo:7.0
    container_name: go_storage_mongo
    restart: unless-stopped
    ports:
      - "27017:27017"

  elasticsearch:
    image: docker.elastic.co/elasticsearch/elasticsearch:8.14.3
    container_name: go_storage_es
    restart: unless-stopped
    environment:
      - discovery.type=single-node
      - xpack.security.enabled=false
      - ES_JAVA_OPTS=-Xms512m -Xmx512m
    ports:
      - "9200:9200"

启动命令:

docker compose up -d

4. 依赖定义

4.1 go.mod

module example.com/go-storage-demo

go 1.22

require (
	github.com/elastic/go-elasticsearch/v8 v8.14.0
	github.com/go-redis/redis/v9 v9.5.1
	github.com/go-sql-driver/mysql v1.8.1
	github.com/jmoiron/sqlx v1.4.0
	go.mongodb.org/mongo-driver v1.15.0
	gorm.io/driver/mysql v1.5.7
	gorm.io/gorm v1.25.10
)

安装依赖:

go mod tidy

5. 数据库迁移(Migration)

很多初学者会直接依赖 GORM 的 AutoMigrate。它在开发期很方便,但在生产环境里,更推荐显式维护 SQL Migration,原因有三点:

  1. 变更历史清晰,可审计
  2. 回滚路径明确
  3. 更容易与 DBA、CI/CD、上线流程集成

5.1 安装迁移工具

这里以 golang-migrate 为例:

go install -tags 'mysql' github.com/golang-migrate/migrate/v4/cmd/migrate@latest

5.2 migrations/000001_init.up.sql

CREATE TABLE users (
    id BIGINT UNSIGNED NOT NULL AUTO_INCREMENT PRIMARY KEY,
    name VARCHAR(128) NOT NULL,
    email VARCHAR(128) NOT NULL,
    created_at DATETIME(3) NOT NULL DEFAULT CURRENT_TIMESTAMP(3),
    updated_at DATETIME(3) NOT NULL DEFAULT CURRENT_TIMESTAMP(3) ON UPDATE CURRENT_TIMESTAMP(3),
    UNIQUE KEY uk_users_email (email)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;

CREATE TABLE orders (
    id BIGINT UNSIGNED NOT NULL AUTO_INCREMENT PRIMARY KEY,
    user_id BIGINT UNSIGNED NOT NULL,
    amount BIGINT NOT NULL,
    status VARCHAR(32) NOT NULL,
    created_at DATETIME(3) NOT NULL DEFAULT CURRENT_TIMESTAMP(3),
    KEY idx_orders_user_id (user_id),
    CONSTRAINT fk_orders_user FOREIGN KEY (user_id) REFERENCES users(id)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;

5.3 migrations/000001_init.down.sql

DROP TABLE IF EXISTS orders;
DROP TABLE IF EXISTS users;

5.4 执行迁移

migrate -path migrations -database "mysql://root:root@tcp(127.0.0.1:3306)/appdb" up

5.5 为什么生产环境要重视 Migration

  • 不要让应用在启动时偷偷改表
  • 表结构变更要能审查、回滚、追踪
  • 每次上线都应把 migration 当成正式发布内容的一部分

6. 数据库连接池配置与优化

Go 的数据库驱动底层都依赖连接池。连接池不是越大越好,而是要结合:

  • 数据库最大连接数
  • 应用实例数
  • 单次 SQL 耗时
  • 峰值并发
  • 数据库 CPU / IO 能力

我们最常调的四个参数如下:

参数 作用 常见建议
SetMaxOpenConns 最大打开连接数 根据数据库容量和实例数控制总量
SetMaxIdleConns 最大空闲连接数 一般设为 MaxOpenConns 的 30%~60%
SetConnMaxLifetime 连接最大生命周期 避免长连接失效、被中间件踢掉
SetConnMaxIdleTime 空闲连接最大存活时间 防止过多无效空闲连接

一个经验值示例:

  • 单实例中等负载服务:MaxOpen=20MaxIdle=10
  • ConnMaxLifetime=30m
  • ConnMaxIdleTime=10m

优化时要注意:

  1. 不是只看应用侧:如果有 10 个应用实例,每个实例 MaxOpenConns=50,总连接数就可能达到 500。
  2. 长 SQL 比小连接池更危险:慢查询会长期占住连接,导致业务看起来像“连接池不够”。
  3. 连接池要配合监控:应定期观察 db.Stats()、慢查询日志、数据库 CPU 和锁等待。

7. 完整可运行示例:main.go

下面给出完整的可运行示例。把它保存为根目录下的 main.go 后,执行 go run . 即可运行。

7.1 main.go

package main

import (
	"bytes"
	"context"
	"database/sql"
	"encoding/json"
	"errors"
	"fmt"
	"log"
	"os"
	"strconv"
	"time"

	"github.com/elastic/go-elasticsearch/v8"
	"github.com/elastic/go-elasticsearch/v8/esapi"
	"github.com/go-redis/redis/v9"
	_ "github.com/go-sql-driver/mysql"
	"github.com/jmoiron/sqlx"
	"go.mongodb.org/mongo-driver/bson"
	"go.mongodb.org/mongo-driver/mongo"
	"go.mongodb.org/mongo-driver/mongo/options"
	"gorm.io/driver/mysql"
	"gorm.io/gorm"
	"gorm.io/gorm/logger"
)

type Config struct {
	MySQLDSN        string
	RedisAddr       string
	RedisPassword   string
	MongoURI        string
	MongoDB         string
	ESAddress       string
	MaxOpenConns    int
	MaxIdleConns    int
	ConnMaxLifetime time.Duration
	ConnMaxIdleTime time.Duration
}

func getenv(key, fallback string) string {
	if v := os.Getenv(key); v != "" {
		return v
	}
	return fallback
}

func getenvInt(key string, fallback int) int {
	if v := os.Getenv(key); v != "" {
		n, err := strconv.Atoi(v)
		if err == nil {
			return n
		}
	}
	return fallback
}

func getenvDuration(key, fallback string) time.Duration {
	v := getenv(key, fallback)
	d, err := time.ParseDuration(v)
	if err != nil {
		panic(fmt.Sprintf("invalid duration for %s: %v", key, err))
	}
	return d
}

func loadConfig() Config {
	return Config{
		MySQLDSN:        getenv("MYSQL_DSN", "root:root@tcp(127.0.0.1:3306)/appdb?charset=utf8mb4&parseTime=True&loc=Local"),
		RedisAddr:       getenv("REDIS_ADDR", "127.0.0.1:6379"),
		RedisPassword:   getenv("REDIS_PASSWORD", ""),
		MongoURI:        getenv("MONGO_URI", "mongodb://127.0.0.1:27017"),
		MongoDB:         getenv("MONGO_DB", "appdb"),
		ESAddress:       getenv("ES_ADDR", "http://127.0.0.1:9200"),
		MaxOpenConns:    getenvInt("DB_MAX_OPEN_CONNS", 20),
		MaxIdleConns:    getenvInt("DB_MAX_IDLE_CONNS", 10),
		ConnMaxLifetime: getenvDuration("DB_CONN_MAX_LIFETIME", "30m"),
		ConnMaxIdleTime: getenvDuration("DB_CONN_MAX_IDLE_TIME", "10m"),
	}
}

type User struct {
	ID        uint64    `gorm:"primaryKey" json:"id"`
	Name      string    `gorm:"size:128;not null" json:"name"`
	Email     string    `gorm:"size:128;not null;uniqueIndex" json:"email"`
	CreatedAt time.Time `json:"created_at"`
	UpdatedAt time.Time `json:"updated_at"`
}

type Order struct {
	ID        uint64    `gorm:"primaryKey" json:"id"`
	UserID    uint64    `gorm:"not null;index" json:"user_id"`
	Amount    int64     `gorm:"not null" json:"amount"`
	Status    string    `gorm:"size:32;not null" json:"status"`
	CreatedAt time.Time `json:"created_at"`
}

type UserProfile struct {
	UserID    uint64         `bson:"user_id" json:"user_id"`
	Tags      []string       `bson:"tags" json:"tags"`
	Ext       map[string]any `bson:"ext" json:"ext"`
	UpdatedAt time.Time      `bson:"updated_at" json:"updated_at"`
}

type UserStats struct {
	UserID      uint64 `db:"user_id" json:"user_id"`
	Name        string `db:"name" json:"name"`
	TotalAmount int64  `db:"total_amount" json:"total_amount"`
	OrderCount  int64  `db:"order_count" json:"order_count"`
}

type App struct {
	cfg         Config
	gormDB      *gorm.DB
	sqlxDB      *sqlx.DB
	redisClient *redis.Client
	mongoClient *mongo.Client
	mongoDB     *mongo.Database
	profileColl *mongo.Collection
	esClient    *elasticsearch.Client
}

func configurePool(db *sql.DB, cfg Config) {
	db.SetMaxOpenConns(cfg.MaxOpenConns)
	db.SetMaxIdleConns(cfg.MaxIdleConns)
	db.SetConnMaxLifetime(cfg.ConnMaxLifetime)
	db.SetConnMaxIdleTime(cfg.ConnMaxIdleTime)
}

func newApp(ctx context.Context, cfg Config) (*App, error) {
	gormDB, err := gorm.Open(mysql.Open(cfg.MySQLDSN), &gorm.Config{
		Logger: logger.Default.LogMode(logger.Warn),
	})
	if err != nil {
		return nil, fmt.Errorf("open gorm mysql failed: %w", err)
	}

	gormSQLDB, err := gormDB.DB()
	if err != nil {
		return nil, fmt.Errorf("get gorm sql db failed: %w", err)
	}
	configurePool(gormSQLDB, cfg)
	if err := gormSQLDB.PingContext(ctx); err != nil {
		return nil, fmt.Errorf("ping gorm mysql failed: %w", err)
	}

	sqlxDB, err := sqlx.Open("mysql", cfg.MySQLDSN)
	if err != nil {
		return nil, fmt.Errorf("open sqlx mysql failed: %w", err)
	}
	configurePool(sqlxDB.DB, cfg)
	if err := sqlxDB.PingContext(ctx); err != nil {
		return nil, fmt.Errorf("ping sqlx mysql failed: %w", err)
	}

	redisClient := redis.NewClient(&redis.Options{
		Addr:     cfg.RedisAddr,
		Password: cfg.RedisPassword,
		DB:       0,
	})
	if err := redisClient.Ping(ctx).Err(); err != nil {
		return nil, fmt.Errorf("ping redis failed: %w", err)
	}

	mongoClient, err := mongo.Connect(ctx, options.Client().ApplyURI(cfg.MongoURI))
	if err != nil {
		return nil, fmt.Errorf("connect mongo failed: %w", err)
	}
	if err := mongoClient.Ping(ctx, nil); err != nil {
		return nil, fmt.Errorf("ping mongo failed: %w", err)
	}
	mongoDB := mongoClient.Database(cfg.MongoDB)
	profileColl := mongoDB.Collection("user_profiles")

	esClient, err := elasticsearch.NewClient(elasticsearch.Config{
		Addresses: []string{cfg.ESAddress},
	})
	if err != nil {
		return nil, fmt.Errorf("create es client failed: %w", err)
	}
	infoRes, err := esClient.Info(esClient.Info.WithContext(ctx))
	if err != nil {
		return nil, fmt.Errorf("ping es failed: %w", err)
	}
	defer infoRes.Body.Close()
	if infoRes.IsError() {
		return nil, fmt.Errorf("es info error: %s", infoRes.String())
	}

	return &App{
		cfg:         cfg,
		gormDB:      gormDB,
		sqlxDB:      sqlxDB,
		redisClient: redisClient,
		mongoClient: mongoClient,
		mongoDB:     mongoDB,
		profileColl: profileColl,
		esClient:    esClient,
	}, nil
}

func (a *App) Close(ctx context.Context) {
	if a.sqlxDB != nil {
		_ = a.sqlxDB.Close()
	}
	if a.gormDB != nil {
		if db, err := a.gormDB.DB(); err == nil {
			_ = db.Close()
		}
	}
	if a.redisClient != nil {
		_ = a.redisClient.Close()
	}
	if a.mongoClient != nil {
		_ = a.mongoClient.Disconnect(ctx)
	}
}

type UserRepository interface {
	CreateUserWithFirstOrder(ctx context.Context, user *User, order *Order) error
	GetUserByID(ctx context.Context, id uint64) (*User, error)
	QueryUserStats(ctx context.Context) ([]UserStats, error)
	SaveUserProfile(ctx context.Context, profile UserProfile) error
	IndexUser(ctx context.Context, user User) error
	SearchUsers(ctx context.Context, keyword string) ([]string, error)
}

type userRepository struct {
	app *App
}

func newUserRepository(app *App) UserRepository {
	return &userRepository{app: app}
}

type UserService struct {
	repo UserRepository
}

func newUserService(repo UserRepository) *UserService {
	return &UserService{repo: repo}
}

func (s *UserService) RegisterUser(ctx context.Context, name, email string, firstOrderAmount int64) (*User, error) {
	user := &User{Name: name, Email: email}
	order := &Order{Amount: firstOrderAmount, Status: "paid"}
	if err := s.repo.CreateUserWithFirstOrder(ctx, user, order); err != nil {
		return nil, err
	}
	return user, nil
}

func (s *UserService) LoadUser(ctx context.Context, id uint64) (*User, error) {
	return s.repo.GetUserByID(ctx, id)
}

func (s *UserService) Stats(ctx context.Context) ([]UserStats, error) {
	return s.repo.QueryUserStats(ctx)
}

func (s *UserService) Search(ctx context.Context, keyword string) ([]string, error) {
	return s.repo.SearchUsers(ctx, keyword)
}

func cacheKeyUser(id uint64) string {
	return fmt.Sprintf("user:%d", id)
}

func (r *userRepository) CreateUserWithFirstOrder(ctx context.Context, user *User, order *Order) error {
	err := r.app.gormDB.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
		if err := tx.Create(user).Error; err != nil {
			return err
		}

		order.UserID = user.ID
		if err := tx.Create(order).Error; err != nil {
			return err
		}

		return nil
	})
	if err != nil {
		return err
	}

	if err := r.app.redisClient.Del(ctx, cacheKeyUser(user.ID)).Err(); err != nil {
		log.Printf("invalidate redis cache failed: %v", err)
	}

	if err := r.SaveUserProfile(ctx, UserProfile{
		UserID: user.ID,
		Tags:   []string{"new-user", "mysql"},
		Ext: map[string]any{
			"first_order_amount": order.Amount,
			"email":              user.Email,
		},
		UpdatedAt: time.Now(),
	}); err != nil {
		log.Printf("sync profile to mongo failed: %v", err)
	}

	if err := r.IndexUser(ctx, *user); err != nil {
		log.Printf("sync user to elasticsearch failed: %v", err)
	}

	return nil
}

func (r *userRepository) GetUserByID(ctx context.Context, id uint64) (*User, error) {
	cached, err := r.app.redisClient.Get(ctx, cacheKeyUser(id)).Result()
	if err == nil {
		var user User
		if unmarshalErr := json.Unmarshal([]byte(cached), &user); unmarshalErr == nil {
			return &user, nil
		}
	} else if err != redis.Nil {
		return nil, err
	}

	var user User
	if err := r.app.gormDB.WithContext(ctx).First(&user, id).Error; err != nil {
		if errors.Is(err, gorm.ErrRecordNotFound) {
			return nil, sql.ErrNoRows
		}
		return nil, err
	}

	payload, err := json.Marshal(user)
	if err == nil {
		if setErr := r.app.redisClient.Set(ctx, cacheKeyUser(id), payload, 5*time.Minute).Err(); setErr != nil {
			log.Printf("set redis cache failed: %v", setErr)
		}
	}

	return &user, nil
}

func (r *userRepository) QueryUserStats(ctx context.Context) ([]UserStats, error) {
	const query = `
SELECT
    u.id AS user_id,
    u.name,
    COALESCE(SUM(o.amount), 0) AS total_amount,
    COUNT(o.id) AS order_count
FROM users u
LEFT JOIN orders o ON o.user_id = u.id
GROUP BY u.id, u.name
ORDER BY total_amount DESC;
`

	stats := make([]UserStats, 0)
	if err := r.app.sqlxDB.SelectContext(ctx, &stats, query); err != nil {
		return nil, err
	}
	return stats, nil
}

func (r *userRepository) SaveUserProfile(ctx context.Context, profile UserProfile) error {
	_, err := r.app.profileColl.UpdateOne(
		ctx,
		bson.M{"user_id": profile.UserID},
		bson.M{"$set": bson.M{
			"tags":       profile.Tags,
			"ext":        profile.Ext,
			"updated_at": profile.UpdatedAt,
		}},
		options.Update().SetUpsert(true),
	)
	return err
}

func (r *userRepository) IndexUser(ctx context.Context, user User) error {
	body, err := json.Marshal(map[string]any{
		"id":         user.ID,
		"name":       user.Name,
		"email":      user.Email,
		"created_at": user.CreatedAt.Format(time.RFC3339),
	})
	if err != nil {
		return err
	}

	req := esapi.IndexRequest{
		Index:      "users",
		DocumentID: strconv.FormatUint(user.ID, 10),
		Body:       bytes.NewReader(body),
		Refresh:    "true",
	}

	res, err := req.Do(ctx, r.app.esClient)
	if err != nil {
		return err
	}
	defer res.Body.Close()

	if res.IsError() {
		return fmt.Errorf("es index error: %s", res.String())
	}
	return nil
}

func (r *userRepository) SearchUsers(ctx context.Context, keyword string) ([]string, error) {
	queryBody, err := json.Marshal(map[string]any{
		"query": map[string]any{
			"multi_match": map[string]any{
				"query":  keyword,
				"fields": []string{"name^2", "email"},
			},
		},
	})
	if err != nil {
		return nil, err
	}

	res, err := r.app.esClient.Search(
		r.app.esClient.Search.WithContext(ctx),
		r.app.esClient.Search.WithIndex("users"),
		r.app.esClient.Search.WithBody(bytes.NewReader(queryBody)),
	)
	if err != nil {
		return nil, err
	}
	defer res.Body.Close()

	if res.IsError() {
		return nil, fmt.Errorf("es search error: %s", res.String())
	}

	var result struct {
		Hits struct {
			Hits []struct {
				Source struct {
					Name string `json:"name"`
				} `json:"_source"`
			} `json:"hits"`
		} `json:"hits"`
	}

	if err := json.NewDecoder(res.Body).Decode(&result); err != nil {
		return nil, err
	}

	names := make([]string, 0, len(result.Hits.Hits))
	for _, hit := range result.Hits.Hits {
		names = append(names, hit.Source.Name)
	}
	return names, nil
}

func mustPrettyJSON(v any) string {
	b, err := json.MarshalIndent(v, "", "  ")
	if err != nil {
		return fmt.Sprintf("marshal failed: %v", err)
	}
	return string(b)
}

func main() {
	ctx := context.Background()
	cfg := loadConfig()

	app, err := newApp(ctx, cfg)
	if err != nil {
		log.Fatalf("init app failed: %v", err)
	}
	defer app.Close(context.Background())

	service := newUserService(newUserRepository(app))

	email := fmt.Sprintf("alice.%d@example.com", time.Now().Unix())
	user, err := service.RegisterUser(ctx, "Alice", email, 19900)
	if err != nil {
		log.Fatalf("register user failed: %v", err)
	}
	fmt.Println("== register user ==")
	fmt.Println(mustPrettyJSON(user))

	firstLoad, err := service.LoadUser(ctx, user.ID)
	if err != nil {
		log.Fatalf("first load user failed: %v", err)
	}
	fmt.Println("== first load user, expected mysql -> redis ==")
	fmt.Println(mustPrettyJSON(firstLoad))

	secondLoad, err := service.LoadUser(ctx, user.ID)
	if err != nil {
		log.Fatalf("second load user failed: %v", err)
	}
	fmt.Println("== second load user, expected redis hit ==")
	fmt.Println(mustPrettyJSON(secondLoad))

	stats, err := service.Stats(ctx)
	if err != nil {
		log.Fatalf("query stats failed: %v", err)
	}
	fmt.Println("== user stats from sqlx ==")
	fmt.Println(mustPrettyJSON(stats))

	searchResult, err := service.Search(ctx, "Alice")
	if err != nil {
		log.Fatalf("search users failed: %v", err)
	}
	fmt.Println("== search result from elasticsearch ==")
	fmt.Println(mustPrettyJSON(searchResult))
}

8. 关系型数据库:GORM / sqlx 实战

8.1 为什么同一个项目里同时使用 GORM 和 sqlx

这不是“重复造轮子”,而是工程上很常见的组合:

  • GORM 负责模型映射、标准 CRUD、事务封装
  • sqlx 负责复杂 SQL、聚合统计、精确控制结果映射

这样做的好处是:

  1. 写业务时不必把所有 CRUD 都写成手工 SQL
  2. 做报表和复杂查询时,不会被 ORM 语法束缚
  3. 团队可以根据场景选择最合适的工具,而不是“一刀切”

8.2 GORM 写入事务示例

核心代码在 CreateUserWithFirstOrder 中:

err := r.app.gormDB.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
	if err := tx.Create(user).Error; err != nil {
		return err
	}

	order.UserID = user.ID
	if err := tx.Create(order).Error; err != nil {
		return err
	}

	return nil
})

这段代码的价值在于:

  • 用户写入成功、订单写入失败时,会自动回滚
  • 不需要手动 Begin / Commit / Rollback
  • 代码更贴近业务含义

8.3 sqlx 复杂查询示例

统计查询通常更适合手写 SQL:

const query = `
SELECT
    u.id AS user_id,
    u.name,
    COALESCE(SUM(o.amount), 0) AS total_amount,
    COUNT(o.id) AS order_count
FROM users u
LEFT JOIN orders o ON o.user_id = u.id
GROUP BY u.id, u.name
ORDER BY total_amount DESC;
`

stats := make([]UserStats, 0)
if err := r.app.sqlxDB.SelectContext(ctx, &stats, query); err != nil {
	return nil, err
}

这里 sqlx 的优势是:

  • SQL 就是最终执行语句,调试简单
  • 查询优化更直观
  • 返回结构体映射足够方便

如果你的系统里报表、排行榜、聚合统计较多,sqlx 往往非常顺手。

9. Redis 客户端使用与缓存策略

9.1 Redis 初始化

示例里使用 go-redis/v9

redisClient := redis.NewClient(&redis.Options{
	Addr:     cfg.RedisAddr,
	Password: cfg.RedisPassword,
	DB:       0,
})

9.2 Cache Aside(旁路缓存)

本文示例采用最经典的 Cache Aside 策略:

  1. 先查 Redis
  2. Redis miss 时查 MySQL
  3. 查到结果后写回 Redis
  4. 更新数据后删除旧缓存

对应代码:

cached, err := r.app.redisClient.Get(ctx, cacheKeyUser(id)).Result()
if err == nil {
	var user User
	if unmarshalErr := json.Unmarshal([]byte(cached), &user); unmarshalErr == nil {
		return &user, nil
	}
}

var user User
if err := r.app.gormDB.WithContext(ctx).First(&user, id).Error; err != nil {
	return nil, err
}

payload, err := json.Marshal(user)
if err == nil {
	_ = r.app.redisClient.Set(ctx, cacheKeyUser(id), payload, 5*time.Minute).Err()
}

更新后删除缓存:

if err := r.app.redisClient.Del(ctx, cacheKeyUser(user.ID)).Err(); err != nil {
	log.Printf("invalidate redis cache failed: %v", err)
}

9.3 常见缓存策略总结

策略 适用场景 实现要点
Cache Aside 最常见业务缓存 读先查缓存,写后删缓存
Read Through 由缓存层统一加载 业务层更简洁,但实现复杂
Write Through 写缓存时同步写库 一致性较强,但链路更长
Write Back 先写缓存、异步落库 性能高,但风险较高

线上项目中,最常见、最稳妥的仍然是 Cache Aside + TTL

9.4 Redis 优化建议

  1. TTL 必须设置:防止冷数据长期占用内存。
  2. 缓存击穿:热点 key 失效时可配合 singleflight、互斥锁。
  3. 缓存穿透:可缓存空值,或使用布隆过滤器。
  4. 缓存雪崩:不同 key 的 TTL 加随机抖动。
  5. 大 key 问题:避免把超大对象直接塞进一个 key。

10. NoSQL:MongoDB、Elasticsearch 集成

10.1 MongoDB 适合放什么数据

MongoDB 的优势是结构灵活,非常适合:

  • 用户画像
  • 行为事件扩展字段
  • 配置快照
  • 半结构化业务文档

本文中,用户画像通过 SaveUserProfile 写入 MongoDB:

_, err := r.app.profileColl.UpdateOne(
	ctx,
	bson.M{"user_id": profile.UserID},
	bson.M{"$set": bson.M{
		"tags":       profile.Tags,
		"ext":        profile.Ext,
		"updated_at": profile.UpdatedAt,
	}},
	options.Update().SetUpsert(true),
)

这种写法的特点:

  • 根据 user_id 做幂等更新
  • 不需要预先定义僵硬的固定列
  • 后续新增字段成本低

10.2 Elasticsearch 适合放什么数据

Elasticsearch 并不适合作为核心事务主库,但非常适合:

  • 全文检索
  • 模糊搜索
  • 多条件筛选
  • 相关性排序
  • 日志与分析检索

本文示例中,用户创建成功后会写入 ES 索引:

req := esapi.IndexRequest{
	Index:      "users",
	DocumentID: strconv.FormatUint(user.ID, 10),
	Body:       bytes.NewReader(body),
	Refresh:    "true",
}

查询时使用 multi_match

queryBody, err := json.Marshal(map[string]any{
	"query": map[string]any{
		"multi_match": map[string]any{
			"query":  keyword,
			"fields": []string{"name^2", "email"},
		},
	},
})

这适合“昵称、邮箱、用户名”等综合搜索场景。

10.3 MySQL、MongoDB、Elasticsearch 如何配合

推荐的思路是:

  • MySQL:保存强一致业务主数据
  • MongoDB:保存灵活扩展信息
  • Elasticsearch:服务于搜索和检索

也就是说:

  • 主数据以 MySQL 为准
  • MongoDB 和 ES 更像是投影、副本、增强视图
  • 不要让核心交易流程依赖 ES 的实时性

11. 数据访问层设计与事务管理

11.1 Repository + Service 分层

示例中采用了比较典型的分层:

  • Repository:直接面向存储
  • Service:组织业务流程
  • main:负责应用组装与启动

Repository 接口如下:

type UserRepository interface {
	CreateUserWithFirstOrder(ctx context.Context, user *User, order *Order) error
	GetUserByID(ctx context.Context, id uint64) (*User, error)
	QueryUserStats(ctx context.Context) ([]UserStats, error)
	SaveUserProfile(ctx context.Context, profile UserProfile) error
	IndexUser(ctx context.Context, user User) error
	SearchUsers(ctx context.Context, keyword string) ([]string, error)
}

Service 层只关心业务动作:

type UserService struct {
	repo UserRepository
}

func (s *UserService) RegisterUser(ctx context.Context, name, email string, firstOrderAmount int64) (*User, error) {
	user := &User{Name: name, Email: email}
	order := &Order{Amount: firstOrderAmount, Status: "paid"}
	if err := s.repo.CreateUserWithFirstOrder(ctx, user, order); err != nil {
		return nil, err
	}
	return user, nil
}

这样设计的好处是:

  1. 业务逻辑与存储实现解耦
  2. 更容易做单元测试
  3. 后续更换存储实现时,影响范围更小

11.2 事务边界应该放在哪里

事务边界通常应放在一个明确的业务动作上,而不是放在单条 SQL 上。

比如“注册用户并创建首单”就是一个完整业务动作,因此事务应包住:

  • 插入 users
  • 插入 orders

而下面这些动作:

  • 写 MongoDB 用户画像
  • 写 Elasticsearch 索引
  • 删 Redis 缓存

通常不应粗暴塞进同一个数据库事务里,因为它们属于跨存储操作。更合理的方式是:

  • 主事务先保证 MySQL 成功提交
  • 再同步或异步刷新缓存、搜索索引、画像文档
  • 必要时引入消息队列做最终一致性

11.3 关于分布式一致性

当系统复杂度上来之后,你会发现:

  • MySQL 负责“真相源”
  • Redis / MongoDB / ES 负责“加速访问”与“增强视图”

所以一致性设计要遵循一个原则:

先保证主库事务正确,再处理外围系统同步。

这比一开始就追求“所有存储绝对同步成功”更现实,也更符合工程实践。

12. 运行方式

按下面步骤即可运行整套示例:

12.1 启动依赖

docker compose up -d

12.2 执行 Migration

migrate -path migrations -database "mysql://root:root@tcp(127.0.0.1:3306)/appdb" up

12.3 拉取依赖并启动程序

go mod tidy
go run .

如果一切正常,你将看到类似输出:

== register user ==
{
  "id": 1,
  "name": "Alice",
  "email": "alice.1710000000@example.com",
  "created_at": "2026-06-08T09:00:00+08:00",
  "updated_at": "2026-06-08T09:00:00+08:00"
}

== first load user, expected mysql -> redis ==
...

== second load user, expected redis hit ==
...

== user stats from sqlx ==
...

== search result from elasticsearch ==
[
  "Alice"
]

13. 生产环境实践建议

13.1 不要把 ORM 当成银弹

GORM 能显著提升开发效率,但:

  • 慢 SQL 分析仍然要回到真实 SQL
  • 大分页、复杂聚合、窗口函数等场景,手写 SQL 往往更可控
  • ORM 适合提升开发效率,不代表所有查询都该 ORM 化

13.2 连接池要配合监控

建议在生产环境里观测以下指标:

  • 打开连接数
  • 空闲连接数
  • 等待连接次数
  • 平均 SQL 延迟
  • 慢查询数量
  • Redis 命中率
  • ES 查询耗时

13.3 Migration 要纳入发布流程

不要让“建表脚本”散落在聊天记录或个人电脑里。正确方式是:

  • 版本化管理 migration
  • 在测试环境先验证
  • 跟随应用版本一起发布
  • 必要时准备回滚脚本

13.4 围绕主库设计一致性

最容易踩坑的一点是:把 Redis、MongoDB、ES 当成和主库同等地位的真相源。更稳妥的做法是:

  • 事务核心数据始终以 MySQL 为准
  • Redis 可失效、可重建
  • ES 可重建索引
  • MongoDB 的扩展文档可根据主库回填

14. 小结

这一章我们围绕「数据库与存储」构建了一套较完整的 Go 实战示例:

  • GORM 完成标准写操作与事务控制
  • sqlx 完成复杂统计查询
  • 连接池参数 管控数据库资源
  • Migration 管理表结构变更
  • Redis 实现 Cache Aside 缓存
  • MongoDB 承载灵活文档数据
  • Elasticsearch 提供搜索能力
  • Repository + Service 做数据访问层解耦

如果你把这套思路真正用到项目里,会发现一个很重要的变化:你不再是“会连数据库”,而是在开始具备存储层架构设计能力


📝 版权声明:本文为原创技术博客,转载请注明出处。

如文章中存在错误或不准确之处,欢迎在评论区指正,感谢您的阅读与支持!

上一篇

Golang项目实践:4.8 综合项目之分布式短链接服务

下一篇

Golang项目实践:4.6 容器化与云原生部署