Параллельное программирование с go-mssqldb

Приложения Go часто используют goroutines для выполнения параллельной работы. Драйвер go-mssqldb и пакет database/sql предназначены для одновременного использования, но необходимо следовать определённым схемам, чтобы избежать утечек соединений, гонок данных и нехватки пула. В этой статье рассматриваются вопросы безопасности в горутинах, пулах работников и плавных режимов отключения.

Безопасность горутин

SQL. DB безопасен для одновременного использования

Экземпляр *sql.DB можно безопасно использовать одновременно из нескольких горутин. Он управляет внутренним пулом соединений и выполняет синхронизацию:

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

Предупреждение

Не создавайте новый sql.Open для каждого запроса или для каждой горутины. Каждый sql.Open вызов создаёт отдельный пул соединений. Создание пулов на каждый запрос тратит ресурсы и может быстро исчерпать ограничения на серверных соединениях.

sql.Rows, sql.Tx и sql.Conn небезопасны для одновременного использования

Эти типы представляют собой одно соединение и должны использоваться из одной горутины за раз:

// 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 для пулов работников

Пакет 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-сервер (или другой маршрутизатор запросов) завершит разряд. Если сначала закрыть пул, у обработчиков на борту появляются ошибки разрыва соединения.

Размер пула соединений для одновременных рабочих нагрузок

Определяйте размер пула соединений с учётом ожидаемого уровня параллелизма вашего приложения, а не общего количества горутин:

Тип приложения Рекомендуется 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, где ограничения соединения зависят от уровня.

Избегайте распространённых ошибок при параллелизации

Не используйте *sql.Rows совместно в нескольких горутинах

Передача *sql.Rows значения нескольким горутинам вызывает гонки данных:

// 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 экземпляр для всего приложения.
Безопасность горутин Не используйте совместно *sql.Rows, *sql.Tx или *sql.Conn между горутинами.
Рабочие пулы Используйте errgroup.SetLimit или семафорный канал для управления параллельностью.
Размер бассейна Установите MaxOpenConns ниже лимита соединения с сервером.
Корректное завершение работы Спустите HTTP-обработчики перед закрытием пула базы данных.
Очистка ресурсов Всегда вызывайте defer rows.Close() и defer tx.Rollback() в той горутине, которая их создала.
Параллельные запросы Выполняйте независимые запросы одновременно, чтобы снизить задержку агрегационных страниц.