Programación concurrente con go-mssqldb

Las aplicaciones de Go suelen usar goroutines para gestionar tareas concurrentes. El go-mssqldb controlador y el database/sql paquete están diseñados para uso simultáneo, pero debes seguir patrones específicos para evitar fugas de conexión, carreras de datos y falta de pool. Este artículo aborda la seguridad en las goroutines, los grupos de trabajadores y los patrones de apagado controlado.

Seguridad de las goroutines

sql.DB es seguro para el uso concurrente

Una *sql.DB instancia es segura para usar desde varias gorutinas simultáneamente. Gestiona un conjunto interno de conexiones y gestiona la sincronización:

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

No crees una nueva sql.Open por petición ni por goroutine. Cada sql.Open llamada crea un pool de conexiones separado. Crear pools por solicitud desperdicia recursos y puede agotar rápidamente los límites de conexión del servidor.

sql.Rows, sql.Tx y sql.Conn no son seguros para su uso concurrente

Estos tipos representan una única conexión y deben utilizarse por una sola goroutine a la vez:

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

Patrón de grupo de trabajadores

Cuando necesites procesar muchos elementos simultáneamente (por ejemplo, actualizar miles de registros), utiliza un grupo de trabajadores de tamaño fijo. Este enfoque limita la concurrencia para evitar el agotamiento de pools y la sobrecarga de servidores:

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
}

Usa errgroup para grupos de trabajadores

El paquete golang.org/x/sync/errgroup simplifica los grupos de trabajo con propagación de errores y cancelación de contexto integradas:

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

Fija el límite errgroup a un valor menor que MaxOpenConns. Si el número de trabajadores es igual al tamaño del grupo, los trabajadores consumen todas las conexiones y no dejan espacio para revisiones de salud u otras consultas.

Consultas paralelas

Ejecuta consultas independientes simultáneamente para reducir la latencia 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
}

Procesamiento por lotes con concurrencia controlada

Para operaciones de grandes lotes (importación de datos, sincronización de registros), combina la concurrencia en lotes para maximizar el rendimiento:

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

Apagado ordenado

Cuando tu aplicación reciba una señal de apagado, drena las operaciones activas de la base de datos antes de cerrar el pool. Si se llama abruptamente a db.Close(), se cancelan las consultas en curso y pueden quedar sesiones huérfanas en el servidor.

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.")
}

Importante

Cierra el pool de bases de datos después de que tu servidor HTTP (u otro router de solicitudes) termine de vaciarse. Si cierras la piscina primero, los manejadores en vuelo tienen errores de conexión rota.

Tamaño del pool de conexiones para cargas de trabajo concurrentes

Dimensiona el conjunto de conexiones en función de la concurrencia esperada de tu aplicación, no del número total de goroutines:

Tipo de aplicación Recomendado MaxOpenConns Justificación
API HTTP, baja concurrencia 10-25 Coincide con el número típico de solicitudes concurrentes.
API HTTP, alta concurrencia 25-50 Más conexiones para manejadores paralelos.
Trabajador de fondo (lote) 5-10 por grupo de trabajadores Cada trabajador necesita su propia conexión.
Mixta (API + trabajos en segundo plano) Suma de API + necesidades de trabajadores Asegúrate de que cada subsistema tenga suficiente margen de mano.
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

Establezca siempre MaxOpenConns. El valor predeterminado (0) es ilimitado. Un pool ilimitado bajo carga puede abrir cientos de conexiones y saturar el servidor, especialmente en Azure SQL, donde los límites de conexión dependen del nivel.

Evitar errores de concurrencia común

No compartas *sql.Rows entre goroutines

Pasar un valor *sql.Rows a varias goroutines provoca carreras de datos:

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

No olvides cerrar filas en bucles

Una fuga *sql.Rows dentro de un circuito expulsa la piscina de conexión:

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

No uses db.Conn a menos que necesites afinidad de conexión.

db.Conn(ctx) fija una conexión específica. Si lo usas innecesariamente, reduces el tamaño efectivo de la piscina:

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

Lista de verificación de concurrencia

Area Recommendation
Compartir el pool Crea una instancia *sql.DB para toda la aplicación.
Seguridad de goroutines No compartas *sql.Rows, *sql.Tx, ni *sql.Conn entre goroutines.
Grupos de trabajadores Usa errgroup.SetLimit o un canal de semáforo para controlar la concurrencia.
Dimensionamiento de piscinas Configura MaxOpenConns por debajo del límite de conexión del servidor.
Apagado ordenado Drena los manejadores HTTP antes de cerrar el pool de bases de datos.
Limpieza de recursos Siempre defer rows.Close() y defer tx.Rollback() en la goroutine que los creó.
Consultas paralelas Ejecuta consultas independientes simultáneamente para reducir la latencia en las páginas de agregación.