Базы данных

Go и PostgreSQL: продвинутые техники работы с pgx, пулом соединений и транзакциями

Ruslan Ismailov Опубликовано 18 мин чтения
G

Введение: почему pgx вместо database/sql

Стандартный интерфейс database/sql в Go универсален, но именно эта универсальность становится узким местом при серьёзной работе с PostgreSQL. Библиотека pgx от jackc предоставляет прямой доступ к протоколу PostgreSQL, минуя абстракции database/sql, что даёт несколько принципиальных преимуществ.

  • Нативная поддержка типов PostgreSQL: JSONB, массивы, UUID, pgtype.Numeric без лишних конвертаций.
  • Поддержка batch-запросов через SendBatch, что радикально сокращает количество round-trips.
  • Встроенный пул соединений pgxpool с детальной настройкой health-check и таймаутов.
  • Поддержка COPY protocol для массовой вставки данных.
  • Прямая работа с prepared statements на уровне протокола.

Для Go-разработчиков, которые строят highload-сервисы на PostgreSQL, переход с database/sql + lib/pq на pgx v5 — одно из самых окупаемых технических решений 2024–2026 годов. Установка:

go get github.com/jackc/pgx/v5
go get github.com/jackc/pgx/v5/pgxpool

Настройка pgxpool: размер пула, таймауты, health check

Пул соединений — центральный элемент производительной работы с PostgreSQL из Go. pgxpool управляет жизненным циклом соединений, их переиспользованием и проверкой работоспособности.

package main

import (
    "context"
    "fmt"
    "log"
    "time"

    "github.com/jackc/pgx/v5/pgxpool"
)

func NewPool(dsn string) (*pgxpool.Pool, error) {
    config, err := pgxpool.ParseConfig(dsn)
    if err != nil {
        return nil, fmt.Errorf("parse config: %w", err)
    }

    // Максимальное количество соединений в пуле
    config.MaxConns = 30
    // Минимальное количество простаивающих соединений
    config.MinConns = 5
    // Максимальное время жизни соединения
    config.MaxConnLifetime = 1 * time.Hour
    // Максимальное время простоя соединения
    config.MaxConnIdleTime = 30 * time.Minute
    // Таймаут ожидания соединения из пула
    config.HealthCheckPeriod = 1 * time.Minute
    // Таймаут на установку соединения
    config.ConnConfig.ConnectTimeout = 5 * time.Second

    // Health check: выполняем SELECT 1 после получения соединения
    config.BeforeAcquire = func(ctx context.Context, conn *pgxpool.Conn) bool {
        return conn.Ping(ctx) == nil
    }

    // Хук после освобождения соединения
    config.AfterRelease = func(conn *pgx.Conn) bool {
        // Сбрасываем prepared statements, если соединение было в ошибке
        return conn.IsClosed() == false
    }

    pool, err := pgxpool.NewWithConfig(context.Background(), config)
    if err != nil {
        return nil, fmt.Errorf("create pool: %w", err)
    }

    return pool, nil
}

Несколько важных правил при настройке pgxpool:

  • MaxConns не должен превышать max_connections PostgreSQL минус соединения для репликации и административных задач. Практическое правило: MaxConns = (количество CPU ядер * 2) + количество дисков.
  • MinConns позволяет держать «тёплые» соединения и избегать latency при всплесках трафика.
  • MaxConnLifetime предотвращает накопление долгоживущих соединений, которые могут держать ресурсы на стороне PostgreSQL.
  • Не злоупотребляйте BeforeAcquire с Ping — это добавляет latency на каждое получение соединения. Включайте только при проблемах с зависшими соединениями.

Работа с транзакциями: явные транзакции, savepoints, обработка ошибок

Транзакции в PostgreSQL через pgx требуют явного управления. Распространённая ошибка — игнорирование отката при панике или ошибке.

func TransferFunds(ctx context.Context, pool *pgxpool.Pool, fromID, toID int64, amount float64) error {
    tx, err := pool.Begin(ctx)
    if err != nil {
        return fmt.Errorf("begin tx: %w", err)
    }
    // Гарантируем откат при выходе из функции с ошибкой
    defer func() {
        if err != nil {
            _ = tx.Rollback(ctx)
        }
    }()

    // Списание
    _, err = tx.Exec(ctx,
        "UPDATE accounts SET balance = balance - $1 WHERE id = $2 AND balance >= $1",
        amount, fromID,
    )
    if err != nil {
        return fmt.Errorf("debit: %w", err)
    }

    // Savepoint перед зачислением
    _, err = tx.Exec(ctx, "SAVEPOINT before_credit")
    if err != nil {
        return fmt.Errorf("savepoint: %w", err)
    }

    // Зачисление
    _, err = tx.Exec(ctx,
        "UPDATE accounts SET balance = balance + $1 WHERE id = $2",
        amount, toID,
    )
    if err != nil {
        // Откат к savepoint, не к началу транзакции
        if rbErr := tx.Exec(ctx, "ROLLBACK TO SAVEPOINT before_credit"); rbErr != nil {
            return fmt.Errorf("rollback to savepoint: %w", rbErr)
        }
        return fmt.Errorf("credit: %w", err)
    }

    // Логируем операцию в той же транзакции
    _, err = tx.Exec(ctx,
        "INSERT INTO audit_log (from_id, to_id, amount, created_at) VALUES ($1, $2, $3, NOW())",
        fromID, toID, amount,
    )
    if err != nil {
        return fmt.Errorf("audit log: %w", err)
    }

    return tx.Commit(ctx)
}

Для уровней изоляции используйте pool.BeginTx с явным указанием уровня:

tx, err := pool.BeginTx(ctx, pgx.TxOptions{
    IsoLevel:   pgx.Serializable,
    AccessMode: pgx.ReadWrite,
})

Паттерн «функция с транзакцией» — удобная обёртка для переиспользования:

func WithTx(ctx context.Context, pool *pgxpool.Pool, fn func(pgx.Tx) error) error {
    tx, err := pool.Begin(ctx)
    if err != nil {
        return err
    }
    defer func() {
        if p := recover(); p != nil {
            _ = tx.Rollback(ctx)
            panic(p)
        } else if err != nil {
            _ = tx.Rollback(ctx)
        } else {
            err = tx.Commit(ctx)
        }
    }()
    err = fn(tx)
    return err
}

Prepared statements и их влияние на производительность

pgx автоматически кэширует prepared statements в режиме extended query protocol. При первом выполнении запрос парсится и планируется на стороне PostgreSQL, последующие вызовы используют кэшированный план.

// Явная подготовка statement
func PrepareStatements(ctx context.Context, conn *pgx.Conn) error {
    _, err := conn.Prepare(ctx, "get_user_by_email",
        "SELECT id, name, email, created_at FROM users WHERE email = $1 AND deleted_at IS NULL",
    )
    if err != nil {
        return fmt.Errorf("prepare get_user_by_email: %w", err)
    }
    return nil
}

// Использование подготовленного statement
func GetUserByEmail(ctx context.Context, conn *pgx.Conn, email string) (*User, error) {
    var u User
    err := conn.QueryRow(ctx, "get_user_by_email", email).Scan(
        &u.ID, &u.Name, &u.Email, &u.CreatedAt,
    )
    if err != nil {
        return nil, err
    }
    return &u, nil
}

Важный нюанс с pgxpool: prepared statements привязаны к конкретному соединению. При использовании пула рекомендуется применять автоматическое кэширование через QueryExecModeSimpleProtocol или настройку StatementCacheCapacity в конфиге соединения:

config.ConnConfig.DefaultQueryExecMode = pgx.QueryExecModeCacheStatement
config.ConnConfig.StatementCacheCapacity = 512

Работа с нестандартными типами PostgreSQL: JSONB, массивы, UUID

Одно из главных преимуществ pgx — нативная работа с типами PostgreSQL без лишних сериализаций.

JSONB

import "github.com/jackc/pgx/v5/pgtype"

type UserSettings struct {
    Theme    string `json:"theme"`
    Language string `json:"language"`
    Notifications bool `json:"notifications"`
}

func SaveUserSettings(ctx context.Context, pool *pgxpool.Pool, userID int64, settings UserSettings) error {
    data, err := json.Marshal(settings)
    if err != nil {
        return err
    }
    _, err = pool.Exec(ctx,
        "UPDATE users SET settings = $1 WHERE id = $2",
        data, userID,
    )
    return err
}

func GetUserSettings(ctx context.Context, pool *pgxpool.Pool, userID int64) (*UserSettings, error) {
    var raw []byte
    err := pool.QueryRow(ctx,
        "SELECT settings FROM users WHERE id = $1",
        userID,
    ).Scan(&raw)
    if err != nil {
        return nil, err
    }
    var s UserSettings
    if err := json.Unmarshal(raw, &s); err != nil {
        return nil, err
    }
    return &s, nil
}

Массивы PostgreSQL

// Вставка массива тегов
func AddTags(ctx context.Context, pool *pgxpool.Pool, articleID int64, tags []string) error {
    _, err := pool.Exec(ctx,
        "UPDATE articles SET tags = $1::text[] WHERE id = $2",
        tags, articleID,
    )
    return err
}

// Чтение массива
func GetTags(ctx context.Context, pool *pgxpool.Pool, articleID int64) ([]string, error) {
    var tags []string
    err := pool.QueryRow(ctx,
        "SELECT tags FROM articles WHERE id = $1",
        articleID,
    ).Scan(&tags)
    return tags, err
}

UUID

import "github.com/google/uuid"

func CreateOrder(ctx context.Context, pool *pgxpool.Pool, userID int64) (uuid.UUID, error) {
    id := uuid.New()
    _, err := pool.Exec(ctx,
        "INSERT INTO orders (id, user_id, status, created_at) VALUES ($1, $2, 'pending', NOW())",
        id, userID,
    )
    if err != nil {
        return uuid.Nil, err
    }
    return id, nil
}

Batch-запросы с pgx: сокращение round-trips

Batch — мощнейший инструмент pgx для highload-сценариев. Вместо N отдельных запросов отправляем один батч и читаем N результатов за один round-trip.

func GetMultipleUsers(ctx context.Context, pool *pgxpool.Pool, ids []int64) ([]*User, error) {
    conn, err := pool.Acquire(ctx)
    if err != nil {
        return nil, err
    }
    defer conn.Release()

    batch := &pgx.Batch{}
    for _, id := range ids {
        batch.Queue("SELECT id, name, email FROM users WHERE id = $1", id)
    }

    results := conn.SendBatch(ctx, batch)
    defer results.Close()

    var users []*User
    for range ids {
        var u User
        err := results.QueryRow().Scan(&u.ID, &u.Name, &u.Email)
        if err != nil {
            return nil, fmt.Errorf("scan user: %w", err)
        }
        users = append(users, &u)
    }

    return users, nil
}

// Batch-вставка
func BulkInsertEvents(ctx context.Context, pool *pgxpool.Pool, events []Event) error {
    conn, err := pool.Acquire(ctx)
    if err != nil {
        return err
    }
    defer conn.Release()

    batch := &pgx.Batch{}
    for _, e := range events {
        batch.Queue(
            "INSERT INTO events (user_id, type, payload, created_at) VALUES ($1, $2, $3, $4)",
            e.UserID, e.Type, e.Payload, e.CreatedAt,
        )
    }

    br := conn.SendBatch(ctx, batch)
    defer br.Close()

    for i := range events {
        if _, err := br.Exec(); err != nil {
            return fmt.Errorf("insert event %d: %w", i, err)
        }
    }
    return nil
}

Для массовой вставки сотен тысяч строк используйте COPY protocol — он ещё быстрее batch-запросов:

func BulkCopyUsers(ctx context.Context, pool *pgxpool.Pool, users []User) error {
    conn, err := pool.Acquire(ctx)
    if err != nil {
        return err
    }
    defer conn.Release()

    rows := make([][]interface{}, len(users))
    for i, u := range users {
        rows[i] = []interface{}{u.Name, u.Email, u.CreatedAt}
    }

    _, err = conn.Conn().CopyFrom(
        ctx,
        pgx.Identifier{"users"},
        []string{"name", "email", "created_at"},
        pgx.CopyFromRows(rows),
    )
    return err
}

Конкурентные обновления: оптимистичные блокировки и SELECT FOR UPDATE

В системах с высокой конкурентностью критически важно правильно обрабатывать параллельные обновления. pgx предоставляет удобный инструментарий для обоих подходов.

SELECT FOR UPDATE (пессимистичная блокировка)

func ReserveProduct(ctx context.Context, pool *pgxpool.Pool, productID int64, quantity int) error {
    return WithTx(ctx, pool, func(tx pgx.Tx) error {
        var stock int
        err := tx.QueryRow(ctx,
            "SELECT stock FROM products WHERE id = $1 FOR UPDATE",
            productID,
        ).Scan(&stock)
        if err != nil {
            return fmt.Errorf("lock product: %w", err)
        }

        if stock < quantity {
            return fmt.Errorf("insufficient stock: have %d, want %d", stock, quantity)
        }

        _, err = tx.Exec(ctx,
            "UPDATE products SET stock = stock - $1 WHERE id = $2",
            quantity, productID,
        )
        return err
    })
}

Оптимистичные блокировки через version/updated_at

type Product struct {
    ID      int64
    Name    string
    Price   float64
    Version int // Счётчик версий
}

func UpdateProductOptimistic(ctx context.Context, pool *pgxpool.Pool, p Product) error {
    result, err := pool.Exec(ctx,
        `UPDATE products 
         SET name = $1, price = $2, version = version + 1 
         WHERE id = $3 AND version = $4`,
        p.Name, p.Price, p.ID, p.Version,
    )
    if err != nil {
        return fmt.Errorf("update: %w", err)
    }

    if result.RowsAffected() == 0 {
        return fmt.Errorf("optimistic lock conflict: product %d was modified concurrently", p.ID)
    }
    return nil
}

// Retry-обёртка для оптимистичных блокировок
func WithOptimisticRetry(maxAttempts int, fn func() error) error {
    for i := 0; i < maxAttempts; i++ {
        err := fn()
        if err == nil {
            return nil
        }
        // Повторяем только при конфликте версий
        if strings.Contains(err.Error(), "optimistic lock conflict") {
            time.Sleep(time.Duration(i*10) * time.Millisecond)
            continue
        }
        return err
    }
    return fmt.Errorf("exceeded max retry attempts (%d)", maxAttempts)
}

Мониторинг пула соединений и диагностика узких мест

pgxpool предоставляет метод Stat() для получения текущего состояния пула. Интегрируйте его с Prometheus или любой другой системой мониторинга.

import (
    "github.com/prometheus/client_golang/prometheus"
    "github.com/prometheus/client_golang/prometheus/promauto"
)

var (
    poolAcquiredConns = promauto.NewGauge(prometheus.GaugeOpts{
        Name: "pgxpool_acquired_connections",
        Help: "Number of currently acquired connections",
    })
    poolIdleConns = promauto.NewGauge(prometheus.GaugeOpts{
        Name: "pgxpool_idle_connections",
        Help: "Number of idle connections in the pool",
    })
    poolTotalConns = promauto.NewGauge(prometheus.GaugeOpts{
        Name: "pgxpool_total_connections",
        Help: "Total number of connections in the pool",
    })
    poolWaitCount = promauto.NewCounter(prometheus.CounterOpts{
        Name: "pgxpool_wait_total",
        Help: "Total number of times waited for a connection",
    })
)

func MonitorPool(pool *pgxpool.Pool, interval time.Duration) {
    ticker := time.NewTicker(interval)
    defer ticker.Stop()
    for range ticker.C {
        stat := pool.Stat()
        poolAcquiredConns.Set(float64(stat.AcquiredConns()))
        poolIdleConns.Set(float64(stat.IdleConns()))
        poolTotalConns.Set(float64(stat.TotalConns()))
        poolWaitCount.Add(float64(stat.EmptyAcquireCount()))
    }
}

Ключевые метрики для анализа производительности PostgreSQL из Go-приложения:

  • AcquiredConns / MaxConns — если близко к 100%, увеличьте пул или оптимизируйте запросы.
  • EmptyAcquireCount — количество ожиданий свободного соединения. Рост этой метрики сигнализирует о недостаточном пуле.
  • MaxConnLifetimeDestroyCount — частое пересоздание соединений из-за lifetime может указывать на слишком короткий MaxConnLifetime.

Для диагностики медленных запросов используйте pg_stat_statements на стороне PostgreSQL и логирование slow queries на стороне Go:

// Middleware для логирования медленных запросов
type LoggingQuerier struct {
    pool      *pgxpool.Pool
    threshold time.Duration
    logger    *slog.Logger
}

func (lq *LoggingQuerier) QueryRow(ctx context.Context, sql string, args ...any) pgx.Row {
    start := time.Now()
    row := lq.pool.QueryRow(ctx, sql, args...)
    elapsed := time.Since(start)
    if elapsed > lq.threshold {
        lq.logger.WarnContext(ctx, "slow query",
            "sql", sql,
            "duration_ms", elapsed.Milliseconds(),
        )
    }
    return row
}

Интеграция с Docker для локальной разработки

Для воспроизводимой локальной разработки с PostgreSQL используйте Docker Compose. Ниже — минимальная конфигурация с настройками производительности.

version: '3.8'
services:
  postgres:
    image: postgres:16-alpine
    environment:
      POSTGRES_DB: myapp
      POSTGRES_USER: myapp
      POSTGRES_PASSWORD: secret
    ports:
      - "5432:5432"
    volumes:
      - postgres_data:/var/lib/postgresql/data
      - ./init.sql:/docker-entrypoint-initdb.d/init.sql
    command: >
      postgres
        -c max_connections=200
        -c shared_buffers=256MB
        -c effective_cache_size=768MB
        -c work_mem=4MB
        -c log_min_duration_statement=100
        -c log_statement=all
    healthcheck:
      test: ["CMD-SHELL", "pg_isready -U myapp"]
      interval: 5s
      timeout: 5s
      retries: 5

volumes:
  postgres_data:

В Go-коде используйте переменные окружения для DSN, чтобы конфигурация работала как локально, так и в production:

func main() {
    dsn := os.Getenv("DATABASE_URL")
    if dsn == "" {
        dsn = "postgres://myapp:secret@localhost:5432/myapp?sslmode=disable"
    }

    pool, err := NewPool(dsn)
    if err != nil {
        log.Fatalf("failed to create pool: %v", err)
    }
    defer pool.Close()

    // Проверяем соединение при старте
    ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
    defer cancel()

    if err := pool.Ping(ctx); err != nil {
        log.Fatalf("failed to ping database: %v", err)
    }
    log.Println("Connected to PostgreSQL")
}

Заключение и best practices

Работа с PostgreSQL из Go через pgx — это не просто выбор драйвера, а целая экосистема инструментов для построения надёжных и производительных сервисов. Подведём итоги в виде ключевых практик:

  1. Всегда используйте pgxpool в production-коде. Прямые соединения через pgx.Connect подходят только для административных задач и тестов.
  2. Настраивайте MaxConns осознанно: ориентируйтесь на возможности PostgreSQL-сервера, а не на произвольные числа.
  3. Defer rollback всегда: паттерн defer tx.Rollback() безопасен — он игнорируется после успешного Commit.
  4. Используйте batch-запросы для операций, которые можно сгруппировать. Экономия на round-trips особенно заметна при высокой latency сети.
  5. Нативные типы вместо строк: UUID, JSONB, массивы — используйте их напрямую, не конвертируйте в string без необходимости.
  6. Мониторьте пул: метрики pgxpool должны быть частью observability вашего сервиса с первого дня.
  7. COPY protocol для bulk insert: если нужно вставить тысячи строк, COPY в 5–10 раз быстрее batch INSERT.
  8. Savepoints для частичных откатов в сложных транзакциях — это стандартный инструмент PostgreSQL, не пренебрегайте им.
  9. Логируйте slow queries: настройте log_min_duration_statement в PostgreSQL и добавьте middleware в Go для корреляции.
  10. Docker для локальной разработки с healthcheck гарантирует, что приложение стартует только после готовности базы.

pgx активно развивается: v5 принёс улучшенный API для работы с типами, более гибкое управление запросами и лучшую производительность. Следите за релизами и обновляйте зависимости — это прямые инвестиции в производительность вашего Go-сервиса на PostgreSQL.

Технологии

Теги

Руслан Исмаилов

Senior Web / Backend разработчик. Senior web/backend разработчик с 9-летним опытом. Стек: PHP, Laravel, PostgreSQL, Redis, Docker, Kubernetes, REST, микросервисы, CI/CD. Подробнее обо мне →