Samtidig programmering med go-mssqldb

Go-applikationer använder ofta goroutines för att hantera samtidiga uppgifter. Drivrutinen go-mssqldb och paketet database/sql är designade för samtidig användning, men du måste följa specifika mönster för att undvika anslutningsläckor, datarace och poolsvält. Den här artikeln behandlar goroutine-säkerhet, arbetarpooler och smidiga avstängningsmönster.

Goroutine-säkerhet

SQL. DB är säker för samtidig användning

En *sql.DB instans är säker att använda från flera goroutines samtidigt. Den hanterar en intern anslutningspool och hanterar synkronisering:

// 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))
}

Varning

Skapa inte en ny sql.Open för varje begäran eller för varje goroutine. Varje sql.Open samtal skapar en separat anslutningspool. Att skapa pooler per förfrågan slösar resurser och kan snabbt uttömma serversidans anslutningsgränser.

SQL. Rader, SQL. Tx och SQL. Conn är osäkra för samtida användning

Dessa typer representerar en enda anslutning och måste användas från en gorutin åt gången:

// 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.
}

Arbetspoolsmönstret

När du behöver behandla många objekt samtidigt (till exempel uppdatera tusentals poster), använd en fast storlek på worker pool. Denna metod begränsar samtidighet för att förhindra poolutmattning och serveröverbelastning:

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
}

Använd errgroup för arbetspooler

Paketet golang.org/x/sync/errgroup förenklar arbetarpooler med inbyggd felspridning och kontextavbokning:

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

Sätt errgruppsgränsen till ett värde lägre än MaxOpenConns. Om antalet arbetare motsvarar poolens storlek förbrukar arbetarna alla anslutningar och lämnar inget utrymme för hälsokontroller eller andra frågor.

Parallella frågor

Kör oberoende frågor samtidigt för att minska total latens:

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
}

Batchbearbetning med kontrollerad samtidighet

För stora batchoperationer (import av data, synkronisering av poster), kombinera batchning med samtidighet för att maximera genomströmningen:

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()
}

Graciös avstängning

När din applikation får en avstängningssignal, töm aktiva databasoperationer innan poolen stängs. Att plötsligt anropa db.Close() avbryter pågående frågor och kan lämna servern med kvarlämnade sessioner utan ägare.

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

Stäng databaspoolen efter att din HTTP-server (eller annan förfrågningsrouter) har tömt färdigt. Om du stänger poolen först får operatörer i flygning fel med trasiga anslutningar.

Storlek på anslutningspoolen för samtidiga arbetsbelastningar

Dimensionera anslutningspoolen utifrån den förväntade samtidigheten i applikationen, inte utifrån det totala antalet goroutines:

Apptyp Rekommenderas MaxOpenConns Motivering
HTTP API, låg samtidighet 10-25 Matchar det typiska antalet samtidiga förfrågningar.
HTTP API, hög samtidighet 25-50 Fler anslutningar för parallella hanterare.
Bakgrundsarbetare (batch) 5–10 per arbetsgrupp Varje arbetare behöver sin egen kontakt.
Blandad (API + bakgrundsjobb) Summan av API + arbetarbehov Se till att varje delsystem har tillräckligt med headroom.
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

Ställ alltid in MaxOpenConns. Standardvärdet (0) är obegränsat. En obegränsad pool under belastning kan öppna hundratals anslutningar och överbelasta servern, särskilt på Azure SQL där anslutningsgränserna är nivåberoende.

Undvik vanliga samtidighetsfel

Dela inte *sql.Rows mellan gorutiner

Att skicka ett *sql.Rows värde till flera gorutiner orsakar datalopp:

// 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.

Glöm inte att stänga rader i loopar

Läckage *sql.Rows inuti en slinga tömmer anslutningspoolen:

// 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()
}

Använd inte db.Conn om du inte behöver anslutningsaffinitet

db.Conn(ctx) fäster en specifik anslutning. Om du använder den i onödan minskar du den effektiva poolstorleken:

// 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")

Checklista för samtidighet

Area Recommendation
Pooldelning Skapa en instans *sql.DB för hela applikationen.
Goroutine-säkerhet Dela inte *sql.Rows, *sql.Tx eller *sql.Conn mellan goroutines.
Arbetarpooler Använd errgroup.SetLimit eller en semaforkanal för att styra samtidigheten.
Poolstorlek Ställ MaxOpenConns in det lägre än serveranslutningsgränsen.
Graciös avstängning Töm HTTP-hanterare innan du stänger databaspoolen.
Resursrensning Alltid defer rows.Close() och defer tx.Rollback() i gorutinen som skapade dem.
Parallella frågor Kör oberoende frågor samtidigt för att minska latens för aggregeringssidor.