Hinweis
Für den Zugriff auf diese Seite ist eine Autorisierung erforderlich. Sie können versuchen, sich anzumelden oder das Verzeichnis zu wechseln.
Für den Zugriff auf diese Seite ist eine Autorisierung erforderlich. Sie können versuchen, das Verzeichnis zu wechseln.
Go-Anwendungen verwenden häufig Goroutinen, um parallele Aufgaben zu verarbeiten. Der go-mssqldb Treiber und das database/sql Paket sind für gleichzeitige Nutzung konzipiert, aber man muss bestimmten Mustern folgen, um Verbindungslecks, Datenrennen und Pool-Hunger zu vermeiden. Dieser Artikel behandelt GoRoutine-Sicherheit, Arbeitspools und elegante Abschaltmuster.
Goroutine-Sicherheit
SQL. DB ist sicher für gleichzeitige Nutzung
Eine Instanz *sql.DB kann aus mehreren GoRoutines gleichzeitig sicher verwendet werden. Es verwaltet einen internen Verbindungspool und übernimmt die Synchronisation:
// 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
Erstelle keine neue sql.Open pro Anfrage oder pro Goroutine. Jeder sql.Open Anruf erstellt einen separaten Verbindungspool. Das Erstellen von Pools pro Anfrage verschwendet Ressourcen und kann die serverseitigen Verbindungsgrenzen schnell ausschöpfen.
sql.Rows, sql.Tx und sql.Conn sind nicht für die gleichzeitige Verwendung geeignet
Diese Typen stellen eine einzelne Verbindung dar und müssen jeweils von einer Goroutine aus verwendet werden:
// 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.
}
Worker Pool-Muster
Wenn Sie viele Artikel gleichzeitig bearbeiten müssen (zum Beispiel Tausende von Datensätzen aktualisieren), verwenden Sie einen Worker Pool mit fester Größe. Dieser Ansatz begrenzt die Nebenläufigkeit, um eine Erschöpfung des Verbindungspools und eine Serverüberlastung zu verhindern:
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
}
Verwenden Sie errgroup für Worker Pools
Das golang.org/x/sync/errgroup-Paket vereinfacht Worker-Pools mit integrierter Fehlerweitergabe und Kontextabbruch:
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
Setze die Errgruppengrenze auf einen Wert niedriger als MaxOpenConns. Wenn die Anzahl der Arbeiter der Poolgröße entspricht, verbrauchen die Arbeiter alle Verbindungen und lassen keinen Raum für Gesundheitschecks oder andere Anfragen.
Parallele Abfragen
Führen Sie unabhängige Abfragen gleichzeitig aus, um die Gesamtlatenz zu reduzieren:
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
}
Batch-Verarbeitung mit kontrollierter Nebenwahl
Für große Stapelverarbeitungen (Daten importieren, Datensätze synchronisieren) kombinieren Sie Stapelverarbeitung mit Nebenläufigkeit, um den Durchsatz zu maximieren:
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()
}
Kontrolliertes Herunterfahren
Wenn Ihre Anwendung ein Abschaltsignal erhält, entleeren Sie aktive Datenbankoperationen, bevor Sie den Pool schließen. Ein abruptes Anrufen db.Close() bricht In-Flight-Anfragen ab und kann den Server mit verwaisten Sitzungen zurücklassen.
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.")
}
Von Bedeutung
Schließe den Datenbankpool, nachdem dein HTTP-Server (oder anderer Request-Router) entladen ist. Wenn du zuerst den Pool schließt, bekommen die In-Flight-Handler fehlerhafte Verbindungsfehler.
Dimensionierung des Verbindungspools für gleichzeitige Workloads
Bemessen Sie den Verbindungspool anhand der erwarteten Nebenläufigkeit Ihrer Anwendung, nicht anhand der Gesamtzahl der Goroutinen:
| Anwendungstyp | Empfohlen MaxOpenConns |
Begründung |
|---|---|---|
| HTTP API, geringe Parallelität | 10-25 | Entspricht der typischen Anzahl gleichzeitiger Anfragen. |
| HTTP API, hohe Parallelität | 25-50 | Mehr Verbindungen für parallele Handler. |
| Hintergrundarbeiter (Batch) | 5-10 pro Mitarbeiterpool | Jeder Arbeiter braucht seine eigene Verbindung. |
| Gemischt (API + Hintergrundjobs) | Summe aus API- und Worker-Anforderungen | Stellen Sie sicher, dass jedes Subsystem genügend Headroom hat. |
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
Immer MaxOpenConns festlegen. Die Standardversion (0) ist unbegrenzt. Ein unbegrenzter Pool unter Last kann Hunderte von Verbindungen öffnen und den Server überlasten, besonders bei Azure SQL, wo die Verbindungsgrenzen stufenabhängig sind.
Vermeiden Sie häufige Fehler bei der Nebenläufigkeit
Teilen Sie *sql.Rows nicht zwischen Goroutinen.
Das Übergeben eines *sql.Rows-Werts an mehrere Goroutinen führt zu Datenrennen:
// 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.
Vergiss nicht, die Zeilen in Schleifen zu schließen
Ein Leck *sql.Rows in einer Schleife erschöpft den Verbindungspool:
// 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()
}
Verwende db.Conn nicht, es sei denn, du benötigst Verbindungsaffinität
db.Conn(ctx) stiftet eine bestimmte Verbindung an. Wenn Sie es unnötig verwenden, verringern Sie die effektive Poolgröße:
// 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")
Nebenläufigkeits-Checkliste
| Area | Recommendation |
|---|---|
| Poolfreigabe | Erstelle eine Instanz *sql.DB für die gesamte Anwendung. |
| Goroutine-Sicherheit | Teile *sql.Rows, *sql.Tx oder *sql.Conn nicht zwischen Goroutinen. |
| Arbeiterpools | Verwenden Sie errgroup.SetLimit oder einen Semaphore-Kanal, um die Nebenläufigkeit zu steuern. |
| Poolgröße | Setze MaxOpenConns niedriger als das Server-Verbindungslimit. |
| Kontrolliertes Herunterfahren | Entleere HTTP-Handler, bevor du den Datenbankpool schließt. |
| Ressourcenbereinigung |
defer rows.Close() und defer tx.Rollback() immer in der Goroutine, die sie erstellt hat. |
| Parallele Abfragen | Führen Sie unabhängige Abfragen gleichzeitig aus, um die Latenz für Aggregationsseiten zu reduzieren. |