Go 应用通常使用 goroutine 来处理并发工作。
go-mssqldb驱动程序和database/sql包设计为并发使用,但必须遵循特定模式以避免连接泄漏、数据冲突和池池枯竭。 本文涵盖 goroutine 的安全性、工作池以及优雅关闭模式。
Goroutine 安全性
SQL。数据库适合并发使用
一个 *sql.DB 实例可安全地被多个 goroutine 同时使用。 它管理内部连接池并处理同步:
// CORRECT: Share a single *sql.DB across all goroutines.
var db *sql.DB
func main() {
var err error
db, err = sql.Open("sqlserver", connString)
if err != nil {
log.Fatal(err)
}
defer db.Close()
http.HandleFunc("/employees", listEmployees) // Each request runs in its own goroutine.
log.Fatal(http.ListenAndServe(":8080", nil))
}
Warning
不要为每个请求或每个流程创建一个新的 sql.Open 。 每个 sql.Open 通话都会创建一个独立的连接池。 按请求创建池会浪费资源,并且可能很快耗尽服务器端连接限制。
sql.Rows、sql.Tx 和 sql.Conn 不能并发使用
这些类型表示单个连接,且同一时间只能由一个 goroutine 使用:
// WRONG: Sharing rows across goroutines causes data races.
rows, _ := db.QueryContext(ctx, "SELECT BusinessEntityID, FirstName + ' ' + LastName FROM Sales.vSalesPerson")
go func() { rows.Next() }() // DATA RACE
go func() { rows.Next() }() // DATA RACE
// CORRECT: Process rows in the goroutine that created them.
rows, _ := db.QueryContext(ctx, "SELECT BusinessEntityID, FirstName + ' ' + LastName FROM Sales.vSalesPerson")
defer rows.Close()
for rows.Next() {
// Process in this goroutine only.
}
工作线程池模式
当你需要同时处理许多项目(例如更新数千条记录)时,可以使用固定大小的工作池。 这种方法限制并发,以防止池耗尽和服务器过载:
import "sync"
func updateEmployeeLocations(ctx context.Context, db *sql.DB, updates []EmployeeUpdate) error {
const maxWorkers = 10
sem := make(chan struct{}, maxWorkers)
var mu sync.Mutex
var firstErr error
var wg sync.WaitGroup
for _, u := range updates {
select {
case <-ctx.Done():
return ctx.Err()
case sem <- struct{}{}: // Acquire a worker slot.
}
wg.Add(1)
go func(u EmployeeUpdate) {
defer wg.Done()
defer func() { <-sem }() // Release the worker slot.
_, err := db.ExecContext(ctx,
"UPDATE HumanResources.Department SET GroupName = @grp WHERE DepartmentID = @id",
sql.Named("grp", u.GroupName),
sql.Named("id", u.Id))
if err != nil {
mu.Lock()
if firstErr == nil {
firstErr = err
}
mu.Unlock()
}
}(u)
}
wg.Wait()
return firstErr
}
使用errgroup管理worker pools。
golang.org/x/sync/errgroup 软件包通过内置的错误传播和上下文取消机制,简化了工作池的实现:
import "golang.org/x/sync/errgroup"
func updateDepartmentGroups(ctx context.Context, db *sql.DB, updates []DepartmentUpdate) error {
g, ctx := errgroup.WithContext(ctx)
g.SetLimit(10) // Maximum concurrent goroutines.
for _, u := range updates {
u := u
g.Go(func() error {
_, err := db.ExecContext(ctx,
"UPDATE HumanResources.Department SET GroupName = @grp WHERE DepartmentID = @id",
sql.Named("grp", u.GroupName),
sql.Named("id", u.Id))
return err
})
}
return g.Wait()
}
Tip
将错误组极限设为小于 MaxOpenConns的值。 如果工作者数量等于池大小,工作者将消耗所有连接,且没有空间进行健康检查或其他查询。
并行查询
并发运行独立查询以减少总延迟:
func getDashboardData(ctx context.Context, db *sql.DB) (*Dashboard, error) {
g, ctx := errgroup.WithContext(ctx)
var orderCount int
var customerCount int
var revenue float64
g.Go(func() error {
return db.QueryRowContext(ctx,
"SELECT COUNT(*) FROM Sales.SalesOrderHeader WHERE OrderDate >= DATEADD(day, -7, GETUTCDATE())").
Scan(&orderCount)
})
g.Go(func() error {
return db.QueryRowContext(ctx,
"SELECT COUNT(DISTINCT CustomerID) FROM Sales.SalesOrderHeader WHERE OrderDate >= DATEADD(day, -7, GETUTCDATE())").
Scan(&customerCount)
})
g.Go(func() error {
return db.QueryRowContext(ctx,
"SELECT ISNULL(SUM(TotalDue), 0) FROM Sales.SalesOrderHeader WHERE OrderDate >= DATEADD(day, -7, GETUTCDATE())").
Scan(&revenue)
})
if err := g.Wait(); err != nil {
return nil, err
}
return &Dashboard{
OrderCount: orderCount,
CustomerCount: customerCount,
WeeklyRevenue: revenue,
}, nil
}
带受控并发的批处理
对于大型批处理操作(导入数据、同步记录),将批处理与并发结合以最大化吞吐量:
func importRecords(ctx context.Context, db *sql.DB, records []Record) error {
const batchSize = 100
const maxWorkers = 5
g, ctx := errgroup.WithContext(ctx)
g.SetLimit(maxWorkers)
for i := 0; i < len(records); i += batchSize {
end := i + batchSize
if end > len(records) {
end = len(records)
}
batch := records[i:end]
g.Go(func() error {
return insertBatch(ctx, db, batch)
})
}
return g.Wait()
}
func insertBatch(ctx context.Context, db *sql.DB, batch []Record) error {
tx, err := db.BeginTx(ctx, nil)
if err != nil {
return err
}
defer tx.Rollback()
stmt, err := tx.Prepare(mssql.CopyIn("Production.ScrapReason", mssql.BulkOptions{}, "Name"))
if err != nil {
return err
}
for _, r := range batch {
if _, err := stmt.Exec(r.Name); err != nil {
return err
}
}
if _, err := stmt.Exec(); err != nil {
return err
}
if err := stmt.Close(); err != nil {
return err
}
return tx.Commit()
}
正常关闭
当你的应用程序收到关闭信号时,先排空当前数据库操作,然后再关闭池。 突然调用 db.Close() 会取消飞行中查询,并且可能导致服务器会话被孤立。
import (
"context"
"database/sql"
"log"
"net/http"
"os"
"os/signal"
"syscall"
"time"
)
func main() {
db, err := sql.Open("sqlserver", connString)
if err != nil {
log.Fatal(err)
}
srv := &http.Server{Addr: ":8080"}
// Run the server in a goroutine.
go func() {
if err := srv.ListenAndServe(); err != http.ErrServerClosed {
log.Fatalf("HTTP server error: %v", err)
}
}()
// Wait for interrupt signal.
quit := make(chan os.Signal, 1)
signal.Notify(quit, syscall.SIGINT, syscall.SIGTERM)
<-quit
log.Println("Shutting down...")
// Give in-flight HTTP requests up to 30 seconds to complete.
shutdownCtx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
defer cancel()
if err := srv.Shutdown(shutdownCtx); err != nil {
log.Printf("HTTP shutdown error: %v", err)
}
// Close the database pool after HTTP handlers have drained.
// This waits for any remaining connections to be returned.
if err := db.Close(); err != nil {
log.Printf("Database close error: %v", err)
}
log.Println("Shutdown complete.")
}
Important
在你的 HTTP 服务器(或其他请求路由器)完成清空请求之后关闭数据库连接池。 如果你先关闭池子,飞行中的处理员会出现连接中断的错误。
针对并发工作负载的连接池大小调整
根据应用程序的预期并发量来设置连接池大小,而不是根据 goroutine 总数:
| 应用程序类型 | 推荐 MaxOpenConns |
理由 |
|---|---|---|
| HTTP API,低并发性 | 10-25 | 与典型并发请求数相匹配。 |
| HTTP API,高并发性 | 25-50 | 为并行处理器提供更多连接。 |
| 后台工作人员(批次) | 每个工人池5-10人 | 每个员工都需要自己的连接。 |
| 混合(API + 后台作业) | API + 工人需求的总和 | 确保每个子系统都留有足够的余量。 |
db.SetMaxOpenConns(25) // Total connections across all goroutines.
db.SetMaxIdleConns(10) // Keep warm connections ready for bursts.
db.SetConnMaxLifetime(5 * time.Minute) // Rotate connections for load balancer compatibility.
Tip
始终设置 MaxOpenConns。 默认0()是无限的。 负载下的无限连接池可能会打开数百个连接,导致服务器不堪重负,尤其是在 Azure SQL 上,连接限制取决于层级。
避免常见的并发错误
不要在 goroutine 之间共享 *sql.Rows
将某个 *sql.Rows 值传递给多个 goroutine 会导致数据竞争:
// WRONG: rows is consumed by two goroutines.
rows, _ := db.QueryContext(ctx, "SELECT BusinessEntityID, FirstName + ' ' + LastName FROM Sales.vSalesPerson")
go processRows(rows)
go processRows(rows) // Race condition.
别忘了在循环中结束各行
环路内部泄漏 *sql.Rows 会耗尽连接池:
// WRONG: Rows leak when the loop starts a new iteration.
for _, location := range locations {
rows, _ := db.QueryContext(ctx,
"SELECT FirstName + ' ' + LastName FROM Sales.vSalesPerson WHERE CountryRegionName = @p1",
sql.Named("p1", location))
for rows.Next() {
// Process...
}
// rows.Close() never called if an error occurs.
}
// CORRECT: Use a helper function with defer.
for _, location := range locations {
if err := processLocation(ctx, db, location); err != nil {
return err
}
}
func processLocation(ctx context.Context, db *sql.DB, location string) error {
rows, err := db.QueryContext(ctx,
"SELECT FirstName + ' ' + LastName FROM Sales.vSalesPerson WHERE CountryRegionName = @p1",
sql.Named("p1", location))
if err != nil {
return err
}
defer rows.Close() // Guaranteed cleanup.
for rows.Next() {
// Process...
}
return rows.Err()
}
除非你需要连接亲和性,否则不要使用 db.Conn
db.Conn(ctx) 固定特定连接。 如果不必要地使用,会减少有效池容量:
// WRONG: Unnecessary pinning.
conn, _ := db.Conn(ctx)
defer conn.Close()
conn.QueryContext(ctx, "SELECT 1") // Use db.QueryContext instead.
// CORRECT: Use db.Conn only for temp tables or session-scoped state.
conn, _ := db.Conn(ctx)
defer conn.Close()
conn.ExecContext(ctx, "CREATE TABLE #Temp (Id INT)")
conn.ExecContext(ctx, "INSERT INTO #Temp VALUES (1)")
conn.QueryContext(ctx, "SELECT * FROM #Temp")
并发检查表
| Area | Recommendation |
|---|---|
| 池共享 | 为整个应用创建一个 *sql.DB 实例。 |
| goroutine 安全性 | 不要在 goroutine 之间共享 *sql.Rows、*sql.Tx 或 *sql.Conn。 |
| 工人池 | 使用 errgroup.SetLimit 或信号量通道来控制并发。 |
| 泳池规模 | 将 MaxOpenConns 设置为低于服务器连接限制的值。 |
| 正常关闭 | 在关闭数据库池之前,先排空HTTP处理程序。 |
| 资源清理 | 始终在创建它们的 goroutine 中 defer rows.Close() 和 defer tx.Rollback()。 |
| 并行查询 | 并发独立查询以减少聚合页面的延迟。 |