Go 應用程式通常使用 goroutine 來處理同時進行的工作。
go-mssqldb 驅動程式和 database/sql 套件設計為可同時使用,但你必須遵循特定的模式,以避免連線外洩、資料競爭和連線池耗盡。 本文將介紹日常運作安全、員工池及優雅的關機模式。
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
不要為每個請求或每個 goroutine 建立新的 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 pool
此 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
將 errgroup 限制設為小於 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()。 |
| 平行查詢 | 同時執行獨立查詢以減少聚合頁面的延遲。 |