Go 数据库:database/sql、连接池与事务

# Go 数据库:database/sql、连接池与事务

本篇目标

看懂 Go 服务从数据库驱动、sql.DB 连接池、Repository 到事务的完整关系,正确处理 Context、Rows、空结果和连接池预算。

理解 sql.DB安全执行 SQL划定事务边界管理连接池

# 先记住一句话

sql.DB 是可被多个 Goroutine 共享的连接池句柄,不是一条连接;查询必须传递 Context 并关闭 sql.Rows,事务中的所有操作必须使用同一个 sql.Tx,不能混回 sql.DB。

# 先认识三个核心类型

sql 是 Go 标准库 database/sql 的包名,DB、Rows 和 Tx 是该包对外提供的类型。代码中更常见到它们的指针形式 *sql.DB、*sql.Rows 和 *sql.Tx。

  • sql.DB:整个数据库连接池的操作入口。sql.Open 返回 *sql.DB,后续查询时由它选择或新建可用连接。它通常在服务启动时创建,被所有请求共享,只在服务退出时关闭。
  • sql.Rows:多行查询结果的读取器。QueryContext 返回 *sql.Rows,代码通过 Next 逐行前进,再用 Scan 取出当前行的各列。它会占用查询相关资源,因此要调用 Close,遍历结束后还要检查 Err。
  • sql.Tx:一次数据库事务的操作入口。BeginTx 返回 *sql.Tx,它会在同一条连接上组织多次读写,最后只能选择 Commit 提交或 Rollback 回滚。事务结束后不能继续使用这个 sql.Tx。

可以先把它们记成:sql.DB 管 “整个连接池”,sql.Rows 管 “这次查询返回的多行数据”,sql.Tx 管 “这次事务里的所有操作”。

# database/sql 与驱动的关系

标准库 database/sql 定义统一 API 和连接池行为,具体 PostgreSQL、MySQL 驱动负责与数据库协议通信。

业务 Repository
└─ database/sql 标准接口
   └─ 数据库驱动
      └─ PostgreSQL / MySQL / SQLite
1
2
3
4

这张图表示 Go 业务代码访问数据库时经过的层次,从上往下理解:

  1. 业务 Repository: 编写 “查询笔记”、“保存用户” 等与业务相关的数据访问代码。
  2. database/sql 标准接口: Go 标准库提供统一的查询、事务和连接池 API,例如 QueryContext、ExecContext 和 BeginTx,但它自己不负责实现 PostgreSQL 或 MySQL 的通信协议。
  3. 数据库驱动: 把 database/sql 的统一调用转换成特定数据库能理解的协议,例如 pgx 可以作为 PostgreSQL 驱动。
  4. 具体数据库: 最终接收请求、执行 SQL 并返回结果的 PostgreSQL、MySQL 或 SQLite。

例如 Repository 执行 db.QueryContext(ctx, query) 时,实际调用过程是:

Repository 调用 QueryContext
→ database/sql 管理连接并调用驱动
→ 数据库驱动发送请求
→ 具体数据库执行 SQL 并返回结果
1
2
3
4

图中的 PostgreSQL、MySQL 和 SQLite 表示项目可以选择其中一种,不是一次查询会依次访问三个数据库。 业务代码使用统一的 database/sql API,驱动负责适配具体数据库。

// 文件位置:internal/database/postgres.go
package database

import (
	"context"      // 为启动连通性检查提供取消和超时控制。
	"database/sql" // 标准数据库 API 和连接池实现。
	"fmt"          // 使用 Errorf 为驱动错误补充操作上下文。
	"time"         // 配置连接空闲时间和最长生命周期。

	// _ 表示空白导入:代码不直接调用这个包,只执行它的初始化逻辑。
	// pgx 会在初始化时向 database/sql 注册名为 "pgx" 的驱动,
	// 因此下面的 sql.Open("pgx", dsn) 才能根据名称找到该驱动。
	_ "github.com/jackc/pgx/v5/stdlib"
)

// OpenPostgres 创建 PostgreSQL 连接池句柄,配置容量,并验证数据库当前可访问。
// ctx 应由启动流程传入且带有超时;dsn 是驱动所需的数据库连接字符串。
func OpenPostgres(ctx context.Context, dsn string) (*sql.DB, error) {
	// sql.Open 创建并配置连接池句柄,通常不会立刻建立数据库连接。
	db, err := sql.Open("pgx", dsn)
	if err != nil {
		return nil, fmt.Errorf("open postgres handle: %w", err)
	}

	// 连接池上限要结合数据库容量与服务最大实例数计算。
	db.SetMaxOpenConns(20)                  // 同一时刻最多存在 20 条打开的底层连接。
	db.SetMaxIdleConns(10)                  // 最多保留 10 条空闲连接供后续请求复用。
	db.SetConnMaxIdleTime(5 * time.Minute)  // 一条连接连续空闲超过 5 分钟后可被回收。
	db.SetConnMaxLifetime(30 * time.Minute) // 一条连接无论是否活跃,最多复用 30 分钟。

	// PingContext 会真正向数据库发起连通性检查,并响应 ctx 的超时或取消。
	if err := db.PingContext(ctx); err != nil {
		_ = db.Close() // 启动校验失败时释放已经创建的句柄。
		return nil, fmt.Errorf("ping postgres: %w", err)
	}

	return db, nil // 返回可在整个应用中共享的连接池句柄。
}
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38

sql.Open 通常只验证参数并创建句柄,不保证数据库已经可连接,因此启动阶段用有超时的 PingContext 做就绪校验。sql.DB 应长期复用,不要每个请求都 Open / Close。

# 后续内容怎样串起来

后面几节不是彼此独立的工具,而是围绕一次数据库调用解决不同问题:

Service 传入 Context(超时、取消)
└─ Repository(编写 SQL,在数据库行与 Go Struct 之间转换)
   └─ DBTX 接口(让 Repository 不绑定一种执行方式)
      ├─ *sql.DB:普通查询,从连接池选择连接 ─┐
      └─ *sql.Tx:事务查询,固定使用同一条连接 ─┴→ 数据库驱动 → PostgreSQL / MySQL / SQLite
1
2
3
4
5
  • Repository 是业务代码的数据访问入口。 先掌握它怎样执行 SQL 和映射结果。
  • 事务是多步写入的一致性边界。 需要 “全部成功或全部回滚” 时,使用同一个 sql.Tx。
  • DBTX 只是一个小接口,不是新的数据库组件。 它抽取 sql.DB 和 sql.Tx 共有的方法,让 Repository 中的查询代码可以在两种场景中复用。
  • Context 不是下一层,而是贯穿整条调用链。 它把请求取消和超时从 Service 传到 Repository,再传给驱动。
  • 连接池属于 sql.DB 的内部资源管理。 它决定一个服务实例最多可以同时使用多少条底层连接。
  • Schema 迁移不在请求调用链中。 它属于发布流程,负责在服务使用新表或新字段之前更新数据库结构。

# 数据访问层:Repository 负责 SQL 与数据映射

// 文件位置:internal/note/repository.go
package note

import (
	"context"      // 把请求取消和截止时间传递给数据库驱动。
	"database/sql" // 使用连接池、查询结果以及 sql.ErrNoRows。
	"errors"       // 沿错误链判断是否属于空结果。
	"fmt"          // 包装数据库错误并补充 Note ID 等上下文。
	"time"         // 映射数据库中的 created_at 时间列。
)

// 哨兵错误是可复用的固定错误值,用于跨层判断 “未找到” 这一类别。
var ErrNotFound = errors.New("note not found")

type Note struct {
	ID        string    // notes.id。
	OwnerID   string    // notes.owner_id。
	Title     string    // notes.title。
	Content   string    // notes.content。
	CreatedAt time.Time // notes.created_at,由驱动转换成 Go 时间值。
}

type Repository struct {
	db *sql.DB // sql.DB 是并发安全的连接池句柄,可以在应用内共享。
}

// NewRepository 注入长期复用的连接池,并返回数据访问对象。
func NewRepository(db *sql.DB) *Repository {
	return &Repository{db: db} // Repository 只保存句柄,不在这里打开或关闭连接。
}

// Get 在指定 Owner 范围内查询一篇 Note,并把数据库行映射成 Note Struct。
func (r *Repository) Get(ctx context.Context, id string, ownerID string) (Note, error) {
	// 反引号创建可跨行的原始字符串;$1、$2 是参数占位符,不要拼接用户输入。
	const query = `
		SELECT id, owner_id, title, content, created_at
		FROM notes
		WHERE id = $1 AND owner_id = $2`

	var result Note // 先创建零值 Struct,Scan 成功后各字段会被数据库列填充。
	// QueryRowContext 的后两个参数依次绑定到 SQL 中的 $1 和 $2。
	// & 取得字段地址,Scan 才能把数据库列写入 result 的各字段。
	err := r.db.QueryRowContext(ctx, query, id, ownerID).Scan(
		&result.ID,
		&result.OwnerID,
		&result.Title,
		&result.Content,
		&result.CreatedAt,
	)
	if errors.Is(err, sql.ErrNoRows) {
		return Note{}, ErrNotFound // 把数据库空结果映射为稳定领域错误。
	}
	if err != nil {
		return Note{}, fmt.Errorf("get note %q: %w", id, err)
	}

	return result, nil // 查询和映射都成功。
}

// List 查询某个 Owner 的最新 Note,并用 limit 限制返回条数。
func (r *Repository) List(ctx context.Context, ownerID string, limit int) ([]Note, error) {
	const query = `
		SELECT id, owner_id, title, content, created_at
		FROM notes
		WHERE owner_id = $1
		ORDER BY created_at DESC
		LIMIT $2`

	// QueryContext 适合多行查询,返回需要遍历和关闭的 *sql.Rows。
	rows, err := r.db.QueryContext(ctx, query, ownerID, limit)
	if err != nil {
		return nil, fmt.Errorf("list notes: %w", err)
	}
	defer rows.Close() // 无论正常结束还是中途返回,都把连接归还给池。

	results := make([]Note, 0) // 初始化为空但非 nil 的结果切片,便于 JSON 编码为空数组。
	// Next 将游标推进到下一行;循环结束后仍要检查 rows.Err()。
	for rows.Next() {
		var item Note // 每一行都使用新的 Struct 接收列值。
		if err := rows.Scan(
			&item.ID,
			&item.OwnerID,
			&item.Title,
			&item.Content,
			&item.CreatedAt,
		); err != nil {
			return nil, fmt.Errorf("scan note: %w", err)
		}
		results = append(results, item) // 当前行映射成功后加入结果集。
	}
	if err := rows.Err(); err != nil {
		return nil, fmt.Errorf("iterate notes: %w", err) // 检查迭代期间的网络错误。
	}

	return results, nil // 即使没有任何行,也返回空切片和 nil。
}
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96

参数必须通过 $1、$2 等占位符绑定,不要用 fmt.Sprintf 拼接用户输入。表名或排序字段不能直接作为普通参数绑定,动态选择时应映射到服务端白名单。

# 多步一致性:事务覆盖一个完整用例

下面的用例同时创建笔记和审计记录,两条写入必须共同成功或共同回滚。

// 文件位置:internal/note/create.go
package note

import (
	"context"      // 事务及事务内 SQL 共享同一请求生命周期。
	"database/sql" // 提供连接池和事务类型。
	"fmt"          // 包装 Begin、Exec 和 Commit 错误。
)

type Creator struct {
	db *sql.DB // 用于开始事务的共享连接池句柄。
}

// NewCreator 注入连接池,返回负责创建 Note 用例的对象。
func NewCreator(db *sql.DB) *Creator {
	return &Creator{db: db}
}

// Create 在同一事务中写入 Note 和审计记录,保证两条数据共同成功或共同回滚。
func (c *Creator) Create(ctx context.Context, note Note) (err error) {
	// (err error) 是命名返回值,下面的 defer 可以观察函数最终返回的错误。
	// 第二个参数 nil 表示使用数据库默认事务隔离级别和读写选项。
	tx, err := c.db.BeginTx(ctx, nil)
	if err != nil {
		return fmt.Errorf("begin create note transaction: %w", err)
	}
	// defer 注册匿名函数,在 Create 返回前执行,用于兜底回滚未提交事务。
	defer func() {
		// Commit 成功后 Rollback 会返回 sql.ErrTxDone;这里忽略它即可。
		_ = tx.Rollback()
	}()

	const insertNote = `
		INSERT INTO notes (id, owner_id, title, content)
		VALUES ($1, $2, $3, $4)`
	// ExecContext 返回 sql.Result 和 error;这里不需要影响行数,因此用 _ 忽略 Result。
	if _, err := tx.ExecContext(
		ctx,
		insertNote,
		note.ID,
		note.OwnerID,
		note.Title,
		note.Content,
	); err != nil {
		return fmt.Errorf("insert note: %w", err)
	}

	const insertAudit = `
		INSERT INTO audit_logs (actor_id, action, resource_id)
		VALUES ($1, $2, $3)`
	// 仍然使用同一个 tx,确保审计记录属于同一事务,而不是从 db 取得另一条连接。
	if _, err := tx.ExecContext(
		ctx,
		insertAudit,
		note.OwnerID,
		"note.created",
		note.ID,
	); err != nil {
		return fmt.Errorf("insert audit log: %w", err)
	}

	if err := tx.Commit(); err != nil {
		return fmt.Errorf("commit create note transaction: %w", err)
	}
	return nil // Commit 成功后事务已经结束,延迟 Rollback 只会得到 sql.ErrTxDone。
}
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66

事务开始后如果某一步又调用 c.db.ExecContext,它可能从池里取得另一条连接,不属于当前事务。事务内必须只使用 tx.QueryContext、tx.ExecContext 等方法,或者把 *sql.Tx 包装进统一查询接口。

# 兼容普通查询与事务:DBTX 复用查询代码

// 文件位置:internal/database/dbtx.go
package database

import (
	"context"
	"database/sql"
)

type DBTX interface {
	// ...any 是可变参数,允许把任意数量的 SQL 参数继续传给驱动。
	// ExecContext 执行不返回数据行的 INSERT、UPDATE 或 DELETE。
	ExecContext(ctx context.Context, query string, args ...any) (sql.Result, error)
	// QueryContext 执行多行查询,调用方必须遍历并关闭 Rows。
	QueryContext(ctx context.Context, query string, args ...any) (*sql.Rows, error)
	// QueryRowContext 表示预期读取一行,真正的查询错误会在 Scan 时返回。
	QueryRowContext(ctx context.Context, query string, args ...any) *sql.Row
}
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17

*sql.DB 和 *sql.Tx 都满足这个小接口,Repository 可以在普通查询时接收 DB,在事务中接收 Tx。接口应由使用方按实际方法定义,不要为了抽象复制整个 database/sql API。

# 超时与取消:Context 贯穿数据库调用

// 文件位置:internal/note/service.go
package note

import (
	"context" // 派生数据库子调用的超时 Context。
	"time"    // 使用 Second 表示 2 秒超时预算。
)

// Get 为单次数据库查询分配 2 秒预算,并继续调用 Repository。
func (s *Service) Get(ctx context.Context, id string, ownerID string) (Note, error) {
	// 子 Context 会在父 ctx 取消或自身 2 秒期限到达时结束,以更早发生者为准。
	dbCtx, cancel := context.WithTimeout(ctx, 2*time.Second)
	defer cancel() // 及时释放计时器资源。

	// 父请求断开或 2 秒到期时,驱动有机会取消数据库操作。
	return s.repository.Get(dbCtx, id, ownerID)
}
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17

超时应有层级预算:HTTP 总超时如果是 3 秒,数据库子调用不能也配置 5 秒。还要区分客户端取消、Context 截止和数据库错误,便于指标与重试判断。

# 并发容量:连接池限制底层连接数

SetMaxOpenConns 限制同时打开的连接数。达到上限后,新查询会等待空闲连接,因此连接池既保护数据库,也可能成为排队点。

数据库连接总预算
├─ 在线 API 实例数 × 单实例 MaxOpenConns
├─ 后台 Worker 实例数 × 单实例 MaxOpenConns
├─ 迁移与运维连接
└─ 故障切换预留
1
2
3
4
5

应观测 db.Stats() 中的 InUse、Idle、WaitCount、WaitDuration,并结合慢查询和数据库负载调整。简单把连接池调大,可能只是把排队从应用转移到数据库。

// 文件位置:internal/observability/database.go
package observability

import (
	"database/sql" // 读取连接池统计数据。
	"log/slog"     // 以键值形式输出结构化日志。
)

// LogDBStats 读取一次连接池快照并输出关键容量与等待指标。
func LogDBStats(db *sql.DB) {
	// Stats 返回连接池当前快照;这些字段可作为结构化日志键值输出。
	stats := db.Stats()
	// Info 的第一个参数是日志消息,后面按 Key、Value 成对传入字段。
	slog.Info(
		"database pool stats",
		"open", stats.OpenConnections,
		"in_use", stats.InUse,
		"idle", stats.Idle,
		"wait_count", stats.WaitCount,
		"wait_duration", stats.WaitDuration,
	)
}
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22

# 结构版本:Schema 迁移与服务启动分开

迁移文件应纳入版本控制,由发布流程中的单独步骤执行。不要让每个 API 副本在启动时同时修改 Schema。

migrations/
├─ 00001_create_notes.up.sql       # 创建 notes 表
├─ 00001_create_notes.down.sql     # 回滚 notes 表
├─ 00002_add_summary.up.sql        # 添加 summary 字段
└─ 00002_add_summary.down.sql      # 回滚 summary 字段
1
2
3
4
5

使用 Goose、Atlas、Migrate 等工具都可以,重点是:版本有顺序、上线前审查 SQL、评估大表锁和数据回填、发布过程只允许一个迁移执行者。

# 高频面试题与回答

1. `sql.DB` 是数据库连接吗?参考答案

不是。它是并发安全的数据库句柄和连接池,按查询需要取得、创建和归还底层连接,通常在应用生命周期内共享。sql.Open 也不保证已连通,启动时应用带超时的 PingContext 验证。

2. 为什么事务内不能混用 `sql.DB`?参考答案

Tx 绑定一条底层连接,而 DB 调用可能从池中取得另一条连接,不属于当前事务。开始事务后,所有需要原子提交的 SQL 都必须通过同一个 Tx 执行,最后明确 Commit 或 Rollback。

3. Rows 为什么必须 Close,还要检查 Rows.Err?参考答案

Rows 持有查询和底层连接相关资源,关闭后连接才能及时归还池。迭代期间仍可能发生网络或驱动错误,Next 返回 false 不只代表正常结束,所以循环后还要检查 Rows.Err。

4. 数据库连接池是不是越大越好?参考答案

不是。连接过少会让请求排队,过多会挤占数据库连接、内存和并发处理能力。应结合服务实例数、数据库上限、查询耗时和峰值并发设置最大打开连接数、最大空闲连接数与连接生命周期,并观察等待次数和耗时再调整。

# 接下来学什么

下一篇学习 Go HTTP、JSON 与中间件,把 Repository、Context 和错误映射接入真实 HTTP 请求链。

# 参考资料

上次更新时间: 2026年09月18日 02:14:27