Skip to content

PostgreSQL 与 database/sql

Phase 04 — Database 涵盖:PostgreSQL · database/sql · Connection Pool


1. 学习目标

完成本知识点后,你应该能够:

  • 安装并连接 PostgreSQL(本地或 Docker)
  • 使用 Go 标准库 database/sql 执行 CRUD 操作
  • 理解 sql.DB 连接池的配置与原理
  • 正确处理 sql.Rowssql.NullString 等类型
  • 使用参数化查询防止 SQL 注入
  • 在数字孪生项目中持久化设备与遥测数据

2. 为什么需要

SQL 语句可以在 psql 里手动执行,但后端服务必须程序化访问数据库:HTTP 请求来了,Handler 需要查设备、写遥测、更新状态。

Go 标准库 database/sql 提供了与具体数据库无关的接口;PostgreSQL 驱动(pgxlib/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

使用 pgxstdlib 适配器,兼容 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()

RollbackCommit 成功后调用无副作用(返回 sql.ErrTxDone)。这种模式保证出错时自动回滚。


6. JavaScript / TypeScript 对比

概念Go database/sqlNode.js
驱动pgx / lib/pqpg
连接池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

关键差异

  1. Go 的 database/sql同步 API,在 Handler 中直接调用即可(每个请求一个 goroutine)
  2. Go 需要显式处理 NULL(sql.Null*),JS 驱动通常直接映射为 null
  3. 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 流程:

  1. 解析 JSON 请求体
  2. 校验 device_id 是否存在(QueryRow
  3. ExecContext 插入 telemetry
  4. 返回 201 Created

连接池配置建议(中小规模数字孪生):

参数原因
MaxOpenConns25避免过多连接
MaxIdleConns10保持热连接
ConnMaxLifetime30min避免陈旧连接

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:16

10. 练习

请独立完成,不要查看答案。代码放在 workspace/phase-04/db-demo/

Level 1 — 基础

练习 1.1:编写程序连接本地 PostgreSQL,Ping 成功后打印 "connected"

练习 1.2:用 ExecContext 插入一台 AGV 设备记录。

练习 1.3:用 QueryRowContext 查询设备名称并打印。

练习 1.4:配置连接池:MaxOpenConns=10MaxIdleConns=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. 学习检查

完成练习后,确认你能回答:

  1. sql.Open 返回的是什么?它等于「已连接数据库」吗?
  2. 为什么必须用 $1 参数占位符而不是字符串拼接?
  3. defer rows.Close()rows.Err() 为什么都重要?
  4. 连接池的 MaxOpenConns 设太大有什么问题?
  5. 事务中 defer tx.Rollback() 的惯用写法是什么?
  6. sql.NullString 什么时候需要?

12. 下一步

已完成下一知识点关系
PostgreSQL, database/sql, Connection PoolRepository封装数据访问逻辑
ORM了解 GORM / sqlc 简化重复代码

建议顺序:先用 database/sql 手写 Repository 理解底层,再了解 ORM 工具的适用场景。


学习导航

上一篇:SQL 基础与表设计 · 对应练习 · 下一篇:Repository 与 ORM