Go 数据库:database/sql、连接池与事务
# Go 数据库:database/sql、连接池与事务
看懂 Go 服务从数据库驱动、sql.DB 连接池、Repository 到事务的完整关系,正确处理 Context、Rows、空结果和连接池预算。
# 先记住一句话
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
2
3
4
这张图表示 Go 业务代码访问数据库时经过的层次,从上往下理解:
- 业务 Repository: 编写 “查询笔记”、“保存用户” 等与业务相关的数据访问代码。
database/sql标准接口: Go 标准库提供统一的查询、事务和连接池 API,例如QueryContext、ExecContext和BeginTx,但它自己不负责实现 PostgreSQL 或 MySQL 的通信协议。- 数据库驱动: 把
database/sql的统一调用转换成特定数据库能理解的协议,例如 pgx 可以作为 PostgreSQL 驱动。 - 具体数据库: 最终接收请求、执行 SQL 并返回结果的 PostgreSQL、MySQL 或 SQLite。
例如 Repository 执行 db.QueryContext(ctx, query) 时,实际调用过程是:
Repository 调用 QueryContext
→ database/sql 管理连接并调用驱动
→ 数据库驱动发送请求
→ 具体数据库执行 SQL 并返回结果
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 // 返回可在整个应用中共享的连接池句柄。
}
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
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。
}
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。
}
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
}
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)
}
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
├─ 迁移与运维连接
└─ 故障切换预留
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,
)
}
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 字段
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 请求链。