Pemrograman bersamaan dengan go-mssqldb

Aplikasi Go biasanya menggunakan goroutine untuk menangani pekerjaan bersamaan. Driver go-mssqldb dan database/sql paket dirancang untuk penggunaan bersamaan, tetapi Anda harus mengikuti pola tertentu untuk menghindari kebocoran koneksi, balapan data, dan kelaparan kolam. Artikel ini membahas keselamatan goroutine, kumpulan pekerja, dan pola shutdown yang anggun.

Keamanan Goroutine

sql. DB aman untuk penggunaan bersamaan

Sebuah instans *sql.DB aman untuk digunakan oleh beberapa goroutine secara bersamaan. Ini mengelola kumpulan koneksi internal dan menangani sinkronisasi:

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

Jangan membuat sql.Open baru untuk setiap permintaan atau setiap goroutine. Setiap sql.Open panggilan membuat kumpulan koneksi terpisah. Membuat kumpulan per permintaan membuang-buang sumber daya dan dapat menghabiskan batas koneksi sisi server dengan cepat.

sql.Rows, sql.Tx, dan sql.Conn tidak aman untuk digunakan secara bersamaan

Jenis ini mewakili satu koneksi dan hanya boleh digunakan oleh satu goroutine pada satu waktu:

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

Pola kumpulan pekerja

Saat Anda perlu memproses banyak item secara bersamaan (misalnya, memperbarui ribuan rekaman), gunakan kumpulan pekerja berukuran tetap. Pendekatan ini membatasi konkurensi untuk mencegah kelelahan kumpulan dan kelebihan server yang berlebihan:

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
}

Menggunakan errgroup untuk kumpulan pekerja

Paket golang.org/x/sync/errgroup menyederhanakan pool pekerja dengan propagasi kesalahan dan pembatalan konteks bawaan:

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

Atur batas errgroup ke nilai yang lebih rendah dari MaxOpenConns. Jika jumlah pekerja sama dengan ukuran kumpulan, pekerja menggunakan semua koneksi dan tidak menyisakan ruang untuk pemeriksaan kesehatan atau kueri lainnya.

Kueri paralel

Jalankan kueri independen secara bersamaan untuk mengurangi latensi total:

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
}

Pemrosesan batch dengan konkurensi terkontrol

Untuk operasi batch berskala besar (mengimpor data, menyinkronkan catatan), gabungkan pemrosesan batch dengan konkurensi untuk memaksimalkan laju pemrosesan:

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

Penonaktifan yang anggun

Saat aplikasi Anda menerima sinyal shutdown, kuras operasi database aktif sebelum menutup kumpulan. Memanggil db.Close() secara mendadak membatalkan kueri yang sedang berjalan dan dapat menyebabkan server memiliki sesi yatim.

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

Tutup kumpulan database setelah server HTTP Anda (atau router permintaan lainnya) selesai menguras. Jika Anda menutup pool terlebih dahulu, handler yang sedang berjalan akan mengalami galat koneksi terputus.

Ukuran kumpulan koneksi untuk beban kerja bersamaan

Sesuaikan ukuran kumpulan koneksi berdasarkan tingkat konkurensi yang diperkirakan untuk aplikasi Anda, bukan berdasarkan jumlah total goroutine:

Jenis aplikasi Direkomendasikan MaxOpenConns Alasan
HTTP API, konkurensi rendah 10-25 Sesuai dengan jumlah permintaan simultan yang umum.
API HTTP, konkurensi tinggi 25-50 Lebih banyak koneksi untuk handler paralel.
Pekerja latar belakang (batch) 5-10 per kumpulan pekerja Setiap pekerja membutuhkan koneksinya sendiri.
Campuran (API + pekerjaan latar belakang) Jumlah kebutuhan API + pekerja Pastikan setiap subsistem memiliki margin kapasitas yang memadai.
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

Selalu atur MaxOpenConns. Nilai default (0) adalah tidak terbatas. Pool tanpa batas saat menerima beban dapat membuka ratusan koneksi dan membuat server kewalahan, terutama pada Azure SQL, yang batas koneksinya bergantung pada tier.

Hindari kesalahan konkurensi umum

Jangan gunakan *sql.Rows bersama antar-goroutine

Meneruskan *sql.Rows nilai ke beberapa goroutine menyebabkan balapan data:

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

Jangan lupa untuk menutup baris dalam loop

Kebocoran *sql.Rows di dalam loop menghabiskan kumpulan koneksi:

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

Jangan gunakan db.Conn kecuali Anda memerlukan afinitas koneksi

db.Conn(ctx) menyematkan koneksi tertentu. Jika Anda menggunakannya secara tidak perlu, Anda mengurangi ukuran kolam efektif:

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

Daftar periksa konkurensi

Area Recommendation
Berbagi kolam renang Buat satu *sql.DB instans untuk seluruh aplikasi.
Keamanan Goroutine Jangan membagikan *sql.Rows, *sql.Tx, atau *sql.Conn antar-goroutine.
Kumpulan pekerja Gunakan errgroup.SetLimit atau saluran semaphore untuk mengontrol konkurensi.
Ukuran kolam renang Tetapkan MaxOpenConns lebih rendah dari batas koneksi server.
Penonaktifan yang anggun Selesaikan semua handler HTTP sebelum menutup pool database.
Pembersihan sumber daya Selalu defer rows.Close() dan defer tx.Rollback() dalam goroutine yang menciptakannya.
Kueri paralel Jalankan kueri independen secara bersamaan untuk mengurangi latensi untuk halaman agregasi.