go-mssqldb ile eşzamanlı programlama

Go uygulamaları genellikle eşzamanlı işleri yapmak için goroutine kullanır. go-mssqldb sürücüsü ve database/sql paketi eşzamanlı kullanım için tasarlanmıştır, ancak bağlantı sızıntılarını, veri yarışlarını ve bağlantı havuzunun tükenmesini önlemek için belirli kullanım kalıplarını izlemelisiniz. Bu makale, rutin güvenliği, çalışan havuzları ve zarif kapanış kalıplarını ele almaktadır.

Goroutine güvenliği

SQL. DB eşzamanlı kullanım için güvenlidir

Bir *sql.DB örneği, birden çok goroutine tarafından eşzamanlı olarak kullanılabilir. Dahili bir bağlantı havuzunu yönetir ve senkronizasyonu yönetir:

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

Her istek veya her goroutine için yeni bir sql.Open oluşturmayın. Her sql.Open çağrı ayrı bir bağlantı havuzu oluşturur. İstek başına havuz oluşturmak kaynak israfı yapar ve sunucu tarafı bağlantı sınırlarını hızla tüketebilir.

SQL. Rows, sql. Tx ve sql. Conn, eşzamanlı kullanım için güvensizdir

Bu türler tek bir bağlantıyı temsil eder ve aynı anda yalnızca bir goroutine tarafından kullanılmalıdır:

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

İşçi havuzu deseni

Birçok öğeyi eşzamanlı olarak işlemeniz gerekiyorsa (örneğin binlerce kaydı güncellemek), sabit boyutta çalışan havuzu kullanın. Bu yaklaşım, havuz tükenmesini ve sunucu aşırı yüklenmesini önlemek için eşzamanlılığı sınırlar:

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
}

Çalışan havuzları için errgroup kullanın

golang.org/x/sync/errgroup paketi, yerleşik hata yayılımı ve bağlam iptali ile iş parçacığı havuzlarını basitleştirir:

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 sınırını MaxOpenConns değerinden daha düşük bir değere ayarlayın. Eğer çalışan sayısı havuz büyüklüğüne eşitse, çalışanlar tüm bağlantıları tüketir ve sağlık kontrolleri veya diğer sorular için yer bırakmaz.

Paralel sorgular

Bağımsız sorguları eşzamanlı çalıştırarak toplam gecikmeyi azaltın:

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
}

Kontrollü eşzamanlılıkla toplu işleme

Büyük toplu işlemler için (verileri içe aktarma, kayıtları eşzamanlama), aktarım hacmini en üst düzeye çıkarmak amacıyla toplu işlemeyi eşzamanlılıkla birleştirin:

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

Sorunsuz kapatma

Uygulamanız kapatma sinyali aldığında, havuzu kapatmadan önce aktif veritabanı işlemlerini boşaltın. db.Close() öğesinin aniden çağrılması, devam eden sorguları iptal edebilir ve sunucuda sahipsiz oturumlar bırakabilir.

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 sunucunuz (veya başka bir istek yönlendiricisi) boşaldıktan sonra veritabanı havuzunu kapatın. Havuzu önce kapatırsanız, uçuş içi yöneticiler bağlantı hataları alır.

Eşzamanlı iş yükleri için bağlantı havuzu boyutlandırması

Bağlantı havuzunu, uygulamanızın beklenen eşzamanlılığına göre boyutlandırın, toplam goroutine sayısına göre değil:

Uygulama türü Önerilen MaxOpenConns Gerekçe
HTTP API, düşük eşzamanlılık 10-25 Tipik eşzamanlı istek sayısına uyuyor.
HTTP API, yüksek eşzamanlılık 25-50 Paralel işleyiciler için daha fazla bağlantı.
Arka plan çalışanı (grup) Her işçi havuzu için 5-10 Her çalışanın kendi bağlantısına ihtiyacı var.
Karışık (API + arka plan işleri) API + işçi ihtiyaçlarının toplamı Her alt sistemin yeterli boşluk alanına sahip olduğundan emin olun.
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

Her zaman MaxOpenConns ayarlayın. Varsayılan (0) sınırsızdır. Yük altında olan sınırsız havuz, özellikle Azure SQL'de bağlantı sınırlarının katmana bağlı olduğu yüzlerce bağlantı açabilir ve sunucuyu bunaltabilir.

Yaygın eşzamanlılık hatalarından kaçının

*sql.Rows'u goroutine'ler arasında paylaşmayın

Birden fazla goroutine'e bir *sql.Rows değer aktarmak veri yarışlarına neden olur:

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

Sıraları döngülerle kapatmayı unutmayın

*sql.Rows'yi bir döngü içinde sızdırmak, bağlantı havuzunu tüketir:

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

Bağlantı bağlılığına ihtiyacınız olmadıkça db.Conn kullanmayın.

db.Conn(ctx) belirli bir bağlantıyı sabitler. Gereksiz yere kullanırsanız, etkili havuz büyüklüğü azalır:

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

Eşzamanlılık kontrol listesi

Area Recommendation
Havuz paylaşımı Tüm uygulama için bir *sql.DB örnek oluşturun.
Goroutine güvenliği *sql.Rows, *sql.Tx veya *sql.Conn öğelerini goroutine'ler arasında paylaşmayın.
İşçi havuzları Eşzamanlılığı denetlemek için errgroup.SetLimit veya bir semafor kanalı kullanın.
Havuz boyutları Sunucu bağlantı sınırından daha düşük ayarlayın MaxOpenConns .
Düzgün kapatma Veritabanı havuzunu kapatmadan önce HTTP yöneticilerini boşaltın.
Kaynak temizleme Her zaman defer rows.Close() ve defer tx.Rollback() işlemlerini, onları oluşturan goroutine içinde gerçekleştirin.
Paralel sorgular Bağımsız sorguları eşzamanlı çalıştırarak toplama sayfalarının gecikmesini azaltın.