外观
PostgreSQL 与 database/sql
Phase 04 — Database 涵盖:
PostgreSQL·database/sql·Connection Pool
1. 学习目标
完成本知识点后,你应该能够:
- 安装并连接 PostgreSQL(本地或 Docker)
- 使用 Go 标准库
database/sql执行 CRUD 操作 - 理解
sql.DB连接池的配置与原理 - 正确处理
sql.Rows、sql.NullString等类型 - 使用参数化查询防止 SQL 注入
- 在数字孪生项目中持久化设备与遥测数据
2. 为什么需要
SQL 语句可以在 psql 里手动执行,但后端服务必须程序化访问数据库:HTTP 请求来了,Handler 需要查设备、写遥测、更新状态。
Go 标准库 database/sql 提供了与具体数据库无关的接口;PostgreSQL 驱动(pgx 或 lib/pq)负责底层通信。sql.DB 内置连接池,避免每次请求都新建 TCP 连接——这对高并发的数字孪生 API 至关重要。
3. 核心概念
3.1 PostgreSQL
开源关系型数据库,特点:
| 特性 | 说明 |
|---|---|
| ACID 事务 | 完整支持 |
| JSONB | 可存半结构化数据 |
| 扩展丰富 | PostGIS(地理)、TimescaleDB(时序) |
| 并发 | MVCC,读写不阻塞 |
3.2 database/sql 架构
应用代码
↓
database/sql(标准接口:Query, Exec, Begin...)
↓
驱动(pgx / lib/pq)
↓
PostgreSQL Server关键类型:
| 类型 | 作用 |
|---|---|
*sql.DB | 连接池,长期持有,并发安全 |
*sql.Tx | 事务 |
*sql.Rows | 多行查询结果,必须 Close() |
*sql.Row | 单行查询结果 |
3.3 Connection Pool(连接池)
sql.Open() 不会立即连接数据库,而是创建连接池。首次 Query / Exec 时才建立连接。
| 方法 | 作用 | 建议值(参考) |
|---|---|---|
SetMaxOpenConns(n) | 最大打开连接数 | CPU 核数 × 2 ~ 4 |
SetMaxIdleConns(n) | 最大空闲连接数 | MaxOpenConns 的一半 |
SetConnMaxLifetime(d) | 连接最大存活时间 | 5 ~ 30 分钟 |
SetConnMaxIdleTime(d) | 空闲连接最大存活 | 5 分钟 |
4. 基础语法
4.1 项目初始化
bash
mkdir -p workspace/phase-04/db-demo
cd workspace/phase-04/db-demo
go mod init db-demo
go get github.com/jackc/pgx/v5/stdlib使用
pgx的stdlib适配器,兼容database/sql接口,也是目前 PostgreSQL 社区推荐驱动。
4.2 连接数据库
go
package main
import (
"context"
"database/sql"
"fmt"
"log"
"time"
_ "github.com/jackc/pgx/v5/stdlib"
)
func main() {
dsn := "postgres://postgres:secret@localhost:5432/digital_twin?sslmode=disable"
db, err := sql.Open("pgx", dsn)
if err != nil {
log.Fatal(err)
}
defer db.Close()
// 连接池配置
db.SetMaxOpenConns(25)
db.SetMaxIdleConns(10)
db.SetConnMaxLifetime(30 * time.Minute)
db.SetConnMaxIdleTime(5 * time.Minute)
ctx, cancel := context.WithTimeout(context.Background(), 3*time.Second)
defer cancel()
if err := db.PingContext(ctx); err != nil {
log.Fatal("database unreachable:", err)
}
fmt.Println("connected to PostgreSQL")
}4.3 CRUD 操作
go
// 插入设备
func insertDevice(ctx context.Context, db *sql.DB, id, name, deviceType string) error {
_, err := db.ExecContext(ctx,
`INSERT INTO devices (id, name, device_type, status)
VALUES ($1, $2, $3, 'idle')
ON CONFLICT (id) DO NOTHING`,
id, name, deviceType,
)
return err
}
// 查询单条
func getDevice(ctx context.Context, db *sql.DB, id string) (string, string, error) {
var name, status string
err := db.QueryRowContext(ctx,
`SELECT name, status FROM devices WHERE id = $1`, id,
).Scan(&name, &status)
if err != nil {
return "", "", err
}
return name, status, nil
}
// 查询多条
func listAGVs(ctx context.Context, db *sql.DB) ([]Device, error) {
rows, err := db.QueryContext(ctx,
`SELECT id, name, status FROM devices WHERE device_type = $1 ORDER BY id`, "agv",
)
if err != nil {
return nil, err
}
defer rows.Close()
var devices []Device
for rows.Next() {
var d Device
if err := rows.Scan(&d.ID, &d.Name, &d.Status); err != nil {
return nil, err
}
devices = append(devices, d)
}
return devices, rows.Err()
}
type Device struct {
ID string
Name string
Status string
}
// 插入遥测
func insertTelemetry(ctx context.Context, db *sql.DB, deviceID string, x, y, speed float64, battery int) error {
_, err := db.ExecContext(ctx,
`INSERT INTO telemetry (device_id, x, y, speed, battery)
VALUES ($1, $2, $3, $4, $5)`,
deviceID, x, y, speed, battery,
)
return err
}4.4 事务
go
func setDeviceMaintenance(ctx context.Context, db *sql.DB, deviceID, reason string) error {
tx, err := db.BeginTx(ctx, nil)
if err != nil {
return err
}
defer tx.Rollback() // 未 Commit 则回滚
if _, err := tx.ExecContext(ctx,
`UPDATE devices SET status = 'maintenance' WHERE id = $1`, deviceID,
); err != nil {
return err
}
if _, err := tx.ExecContext(ctx,
`INSERT INTO maintenance_logs (device_id, reason) VALUES ($1, $2)`, deviceID, reason,
); err != nil {
return err
}
return tx.Commit()
}4.5 处理 NULL 值
go
import "database/sql"
type TelemetryRow struct {
DeviceID string
X, Y float64
Speed sql.NullFloat64 // 可能为 NULL
Battery sql.NullInt32
RecordedAt time.Time
}
func scanTelemetry(rows *sql.Rows) (TelemetryRow, error) {
var t TelemetryRow
err := rows.Scan(&t.DeviceID, &t.X, &t.Y, &t.Speed, &t.Battery, &t.RecordedAt)
return t, err
}
// 使用时
if t.Speed.Valid {
fmt.Println(t.Speed.Float64)
}5. 代码解析
go
db, err := sql.Open("pgx", dsn)sql.Open 创建的是连接池而非单个连接。即使 Open 成功,也要 Ping 验证数据库可达。
go
_, err := db.ExecContext(ctx, `INSERT ... VALUES ($1, $2)`, id, name)- 使用
$1, $2参数占位符,驱动负责转义,防止 SQL 注入 - PostgreSQL 用
$n;MySQL 用? - 务必传
context,支持超时和取消
go
defer rows.Close()QueryContext 返回的 Rows 必须关闭,否则连接泄漏。在 for rows.Next() 之后检查 rows.Err()。
go
defer tx.Rollback()
return tx.Commit()Rollback 在 Commit 成功后调用无副作用(返回 sql.ErrTxDone)。这种模式保证出错时自动回滚。
6. JavaScript / TypeScript 对比
| 概念 | Go database/sql | Node.js |
|---|---|---|
| 驱动 | pgx / lib/pq | pg 包 |
| 连接池 | sql.DB 内置 | pg.Pool |
| 查询 | db.QueryContext(ctx, sql, args...) | pool.query(sql, params) |
| 参数占位 | $1, $2 | $1, $2(pg 相同) |
| 异步模型 | 同步阻塞(配合 goroutine) | Promise / async-await |
| NULL 处理 | sql.NullString 等 | 直接 null |
关键差异:
- Go 的
database/sql是同步 API,在 Handler 中直接调用即可(每个请求一个 goroutine) - Go 需要显式处理 NULL(
sql.Null*),JS 驱动通常直接映射为null - Go 的
defer rows.Close()是资源管理习惯,类似 JS 的try/finally
7. 常见错误
错误 1:不关闭 Rows
go
// ❌ 连接泄漏
rows, _ := db.QueryContext(ctx, "SELECT ...")
for rows.Next() { ... }
// ✅
rows, err := db.QueryContext(ctx, "SELECT ...")
if err != nil { return err }
defer rows.Close()错误 2:字符串拼接 SQL
go
// ❌ SQL 注入风险
query := fmt.Sprintf("SELECT * FROM devices WHERE id = '%s'", userInput)
// ✅ 参数化查询
db.QueryRowContext(ctx, "SELECT * FROM devices WHERE id = $1", userInput)错误 3:不检查 rows.Err()
go
for rows.Next() {
rows.Scan(&id, &name)
}
// ❌ 遗漏迭代过程中的错误
return devices, nil
// ✅
return devices, rows.Err()错误 4:连接池配置不当
go
// ❌ 默认 MaxOpenConns 无限制,可能压垮数据库
db, _ := sql.Open("pgx", dsn)
// ✅ 根据负载设置
db.SetMaxOpenConns(25)
db.SetMaxIdleConns(10)错误 5:sql.Open 后不做 Ping
go
db, err := sql.Open("pgx", dsn) // 可能 DSN 错误但不报错
// 直到第一次 Query 才发现连不上
if err := db.PingContext(ctx); err != nil {
log.Fatal(err)
}8. 实际应用
数字孪生 — 设备数据持久化 API
REST API 接收 AGV 上报,写入 PostgreSQL:
POST /api/v1/telemetry
{
"device_id": "AGV-001",
"x": 12.5, "y": 3.8,
"speed": 1.2, "battery": 85
}Handler 流程:
- 解析 JSON 请求体
- 校验
device_id是否存在(QueryRow) ExecContext插入telemetry- 返回 201 Created
连接池配置建议(中小规模数字孪生):
| 参数 | 值 | 原因 |
|---|---|---|
| MaxOpenConns | 25 | 避免过多连接 |
| MaxIdleConns | 10 | 保持热连接 |
| ConnMaxLifetime | 30min | 避免陈旧连接 |
9. 深入理解
9.1 sql.DB vs sql.Conn
sql.DB:连接池,绝大多数场景使用sql.Conn:从池中取出的单个连接,用于需要会话级状态的场景(如LISTEN)
9.2 Prepared Statements
go
stmt, err := db.PrepareContext(ctx, "SELECT name FROM devices WHERE id = $1")
defer stmt.Close()
stmt.QueryRowContext(ctx, id).Scan(&name)重复执行的语句预编译可提升性能,但 database/sql 会自动缓存,通常直接 ExecContext 即可。
9.3 Context 传递
始终传递 context.Context:
- HTTP 请求取消时,数据库查询也会取消
- 可设超时:
context.WithTimeout(ctx, 5*time.Second)
9.4 Docker 快速启动 PostgreSQL
bash
docker run -d --name pg-digital-twin \
-e POSTGRES_PASSWORD=secret \
-e POSTGRES_DB=digital_twin \
-p 5432:5432 \
postgres:1610. 练习
请独立完成,不要查看答案。代码放在
workspace/phase-04/db-demo/。
Level 1 — 基础
练习 1.1:编写程序连接本地 PostgreSQL,Ping 成功后打印 "connected"。
练习 1.2:用 ExecContext 插入一台 AGV 设备记录。
练习 1.3:用 QueryRowContext 查询设备名称并打印。
练习 1.4:配置连接池:MaxOpenConns=10,MaxIdleConns=5,打印当前 db.Stats()。
Level 2 — 应用
练习 2.1:实现 ListDevices(ctx, db, deviceType string) ([]Device, error),返回指定类型的所有设备。
练习 2.2:实现 InsertTelemetry(ctx, db, ...) 插入遥测数据。
练习 2.3:实现事务函数:更新设备状态 + 插入告警,任一步失败则回滚。
练习 2.4:处理 sql.ErrNoRows:设备不存在时返回自定义错误信息。
Level 3 — 综合
练习 3.1:编写 HTTP Handler POST /api/v1/telemetry,解析 JSON 并写入数据库。
练习 3.2:编写 GET /api/v1/devices/{id}/latest-telemetry,返回设备最新一条遥测。
练习 3.3:为所有数据库操作添加 5 秒超时 Context。
Level 4 — 项目实践
练习 4.1:搭建完整 db-demo 项目:连接池配置、schema 迁移 SQL、Device CRUD + Telemetry 写入 API,用 curl 或 Postman 验证。
11. 学习检查
完成练习后,确认你能回答:
sql.Open返回的是什么?它等于「已连接数据库」吗?- 为什么必须用
$1参数占位符而不是字符串拼接? defer rows.Close()和rows.Err()为什么都重要?- 连接池的
MaxOpenConns设太大有什么问题? - 事务中
defer tx.Rollback()的惯用写法是什么? sql.NullString什么时候需要?
12. 下一步
| 已完成 | 下一知识点 | 关系 |
|---|---|---|
| PostgreSQL, database/sql, Connection Pool | Repository | 封装数据访问逻辑 |
| ORM | 了解 GORM / sqlc 简化重复代码 |
建议顺序:先用 database/sql 手写 Repository 理解底层,再了解 ORM 工具的适用场景。