Backend-разработка

Построение надёжного Data Pipeline с Go, PostgreSQL и Elasticsearch: от сырых данных до поиска в реальном времени

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

Введение: задача синхронизации данных

Современные приложения всё чаще комбинируют реляционные базы данных с полнотекстовыми поисковыми движками. PostgreSQL надёжно хранит транзакционные данные, но его возможности полнотекстового поиска ограничены по сравнению с Elasticsearch. Задача синхронизации PostgreSQL и Elasticsearch — один из ключевых архитектурных вызовов для backend-разработчиков в 2026 году.

Прямолинейный подход «писать одновременно в обе системы из приложения» ненадёжен: при сбое одной из операций данные расходятся. Правильное решение — выделенный data pipeline, который отслеживает изменения в PostgreSQL и атомарно реплицирует их в Elasticsearch. Именно такой pipeline мы построим в этой статье с использованием Go.

Архитектура pipeline: polling, triggers и CDC

Существуют три основных подхода к синхронизации данных между PostgreSQL и Elasticsearch:

  • Polling — периодический запрос таблиц по полю updated_at. Прост в реализации, но имеет задержку и пропускает удаления.
  • Database Triggers — триггеры PostgreSQL записывают изменения в outbox-таблицу. Надёжнее polling, но увеличивает нагрузку на БД.
  • Change Data Capture (CDC) — чтение журнала WAL (Write-Ahead Log) PostgreSQL. Минимальная нагрузка, поддержка всех операций (INSERT, UPDATE, DELETE), задержка в миллисекунды.

Для production-систем оптимальный выбор — CDC через логическую репликацию PostgreSQL. Это единственный подход, гарантирующий полноту данных без дополнительной нагрузки на приложение. Именно его мы реализуем на Go.

Change Data Capture с PostgreSQL: настройка логической репликации

Для работы с CDC необходимо включить логическую репликацию в PostgreSQL. Отредактируйте postgresql.conf:

wal_level = logical
max_replication_slots = 4
max_wal_senders = 4

Создайте слот репликации и publication для нужных таблиц:

-- Создание publication
CREATE PUBLICATION products_pub FOR TABLE products, categories;

-- Создание слота логической репликации
SELECT pg_create_logical_replication_slot('pipeline_slot', 'pgoutput');

Для чтения WAL-событий в Go мы используем библиотеку pglogrepl — чистый Go-клиент для протокола логической репликации PostgreSQL. Альтернатива — Debezium (JVM-решение), которое удобно, если инфраструктура уже использует Kafka. В нашем случае Go + pglogrepl даёт минимальные зависимости и максимальный контроль.

Реализация pipeline на Go

Структура проекта

Организуем код по принципу чистой архитектуры:

pipeline/
├── cmd/
│   └── pipeline/
│       └── main.go
├── internal/
│   ├── cdc/
│   │   ├── reader.go       # Чтение WAL-событий
│   │   └── decoder.go      # Декодирование pgoutput
│   ├── transform/
│   │   └── mapper.go       # Маппинг PG → ES документы
│   ├── indexer/
│   │   └── elasticsearch.go # Батчевая индексация
│   └── metrics/
│       └── prometheus.go   # Метрики
├── docker-compose.yml
└── Dockerfile

Чтение WAL-событий

Основной worker подключается к PostgreSQL через протокол репликации и читает поток изменений:

package cdc

import (
    "context"
    "fmt"
    "time"

    "github.com/jackc/pglogrepl"
    "github.com/jackc/pgx/v5/pgconn"
    "github.com/jackc/pgx/v5/pgproto3"
)

type WALReader struct {
    conn        *pgconn.PgConn
    slotName    string
    publication string
    lsn         pglogrepl.LSN
}

func NewWALReader(dsn, slotName, publication string) (*WALReader, error) {
    conn, err := pgconn.Connect(context.Background(), dsn+" replication=database")
    if err != nil {
        return nil, fmt.Errorf("connect replication: %w", err)
    }
    return &WALReader{
        conn:        conn,
        slotName:    slotName,
        publication: publication,
    }, nil
}

func (r *WALReader) Start(ctx context.Context, events chan<- *WALEvent) error {
    opts := pglogrepl.StartReplicationOptions{
        PluginArgs: []string{
            "proto_version '1'",
            fmt.Sprintf("publication_names '%s'", r.publication),
        },
    }
    if err := pglogrepl.StartReplication(ctx, r.conn, r.slotName, r.lsn, opts); err != nil {
        return fmt.Errorf("start replication: %w", err)
    }

    standbyDeadline := time.Now().Add(10 * time.Second)
    for {
        if time.Now().After(standbyDeadline) {
            if err := pglogrepl.SendStandbyStatusUpdate(ctx, r.conn,
                pglogrepl.StandbyStatusUpdate{WALWritePosition: r.lsn}); err != nil {
                return fmt.Errorf("standby status: %w", err)
            }
            standbyDeadline = time.Now().Add(10 * time.Second)
        }

        ctx2, cancel := context.WithDeadline(ctx, standbyDeadline)
        msg, err := r.conn.ReceiveMessage(ctx2)
        cancel()
        if err != nil {
            if pgconn.Timeout(err) {
                continue
            }
            return fmt.Errorf("receive message: %w", err)
        }

        switch m := msg.(type) {
        case *pgproto3.CopyData:
            if m.Data[0] == pglogrepl.XLogDataByteID {
                xld, err := pglogrepl.ParseXLogData(m.Data[1:])
                if err != nil {
                    return fmt.Errorf("parse xlog: %w", err)
                }
                event, err := DecodeWALData(xld.WALData)
                if err == nil && event != nil {
                    events <- event
                    r.lsn = xld.WALStart + pglogrepl.LSN(len(xld.WALData))
                }
            }
        }
    }
}

Батчевая индексация в Elasticsearch

Для эффективной индексации используем Bulk API Elasticsearch. Накапливаем события в буфере и сбрасываем по размеру или таймауту:

package indexer

import (
    "bytes"
    "context"
    "encoding/json"
    "fmt"
    "time"

    "github.com/elastic/go-elasticsearch/v8"
    "github.com/elastic/go-elasticsearch/v8/esapi"
)

type BulkIndexer struct {
    client    *elasticsearch.Client
    index     string
    batchSize int
    flushInterval time.Duration
    buf       []BulkAction
}

type BulkAction struct {
    ID     string
    Doc    map[string]interface{}
    Delete bool
}

func (bi *BulkIndexer) Run(ctx context.Context, actions <-chan BulkAction) error {
    ticker := time.NewTicker(bi.flushInterval)
    defer ticker.Stop()

    for {
        select {
        case action, ok := <-actions:
            if !ok {
                return bi.flush(ctx)
            }
            bi.buf = append(bi.buf, action)
            if len(bi.buf) >= bi.batchSize {
                if err := bi.flush(ctx); err != nil {
                    return err
                }
            }
        case <-ticker.C:
            if len(bi.buf) > 0 {
                if err := bi.flush(ctx); err != nil {
                    return err
                }
            }
        case <-ctx.Done():
            return ctx.Err()
        }
    }
}

func (bi *BulkIndexer) flush(ctx context.Context) error {
    if len(bi.buf) == 0 {
        return nil
    }
    var body bytes.Buffer
    for _, action := range bi.buf {
        if action.Delete {
            meta := map[string]interface{}{"delete": map[string]interface{}{"_index": bi.index, "_id": action.ID}}
            line, _ := json.Marshal(meta)
            body.Write(line)
            body.WriteByte('\n')
        } else {
            meta := map[string]interface{}{"index": map[string]interface{}{"_index": bi.index, "_id": action.ID}}
            line, _ := json.Marshal(meta)
            body.Write(line)
            body.WriteByte('\n')
            doc, _ := json.Marshal(action.Doc)
            body.Write(doc)
            body.WriteByte('\n')
        }
    }
    req := esapi.BulkRequest{Body: &body}
    res, err := req.Do(ctx, bi.client)
    if err != nil {
        return fmt.Errorf("bulk request: %w", err)
    }
    defer res.Body.Close()
    if res.IsError() {
        return fmt.Errorf("bulk response error: %s", res.Status())
    }
    bi.buf = bi.buf[:0]
    return nil
}

Гарантии доставки: at-least-once и идемпотентность

WAL-репликация в PostgreSQL гарантирует at-least-once доставку: при перезапуске pipeline события могут быть прочитаны повторно. Elasticsearch Bulk API с явным _id документа обеспечивает идемпотентность — повторная индексация того же документа безопасна.

Критически важно сохранять LSN (Log Sequence Number) после успешной отправки батча в Elasticsearch, а не после чтения из WAL. Храните LSN в отдельной таблице PostgreSQL или Redis:

// Сохранение LSN после успешного flush
func saveLSN(ctx context.Context, db *pgxpool.Pool, slotName string, lsn pglogrepl.LSN) error {
    _, err := db.Exec(ctx,
        `INSERT INTO pipeline_checkpoints (slot_name, lsn, updated_at)
         VALUES ($1, $2, now())
         ON CONFLICT (slot_name) DO UPDATE SET lsn = $2, updated_at = now()`,
        slotName, lsn.String(),
    )
    return err
}

Обработка ошибок при flush должна включать экспоненциальный backoff с jitter, чтобы не перегружать Elasticsearch при временных сбоях:

func retryFlush(ctx context.Context, fn func() error, maxRetries int) error {
    backoff := 100 * time.Millisecond
    for i := 0; i < maxRetries; i++ {
        if err := fn(); err != nil {
            if i == maxRetries-1 {
                return err
            }
            jitter := time.Duration(rand.Int63n(int64(backoff)))
            time.Sleep(backoff + jitter)
            backoff *= 2
            continue
        }
        return nil
    }
    return nil
}

Трансформация данных: маппинг PostgreSQL → Elasticsearch

Структуры данных в PostgreSQL и документы Elasticsearch часто различаются: нормализованные таблицы нужно денормализовать, типы — преобразовать, поля — отфильтровать или обогатить.

package transform

import "time"

// PostgreSQL row
type ProductRow struct {
    ID          int64
    Name        string
    Description string
    Price       float64
    CategoryID  int64
    Category    string
    Tags        []string
    CreatedAt   time.Time
    UpdatedAt   time.Time
    Deleted     bool
}

// Elasticsearch document
type ProductDoc struct {
    ID          string    `json:"id"`
    Name        string    `json:"name"`
    Description string    `json:"description"`
    Price       float64   `json:"price"`
    Category    string    `json:"category"`
    Tags        []string  `json:"tags"`
    UpdatedAt   time.Time `json:"updated_at"`
}

func MapProductToDoc(row ProductRow) ProductDoc {
    return ProductDoc{
        ID:          fmt.Sprintf("%d", row.ID),
        Name:        row.Name,
        Description: row.Description,
        Price:       row.Price,
        Category:    row.Category,
        Tags:        row.Tags,
        UpdatedAt:   row.UpdatedAt,
    }
}

func DocToMap(doc ProductDoc) map[string]interface{} {
    return map[string]interface{}{
        "id":          doc.ID,
        "name":        doc.Name,
        "description": doc.Description,
        "price":       doc.Price,
        "category":    doc.Category,
        "tags":        doc.Tags,
        "updated_at":  doc.UpdatedAt,
    }
}

Обратите внимание: при получении DELETE-события из WAL pipeline должен вызвать удаление документа в Elasticsearch по _id, не запрашивая данные из PostgreSQL (строка уже удалена).

Мониторинг pipeline: метрики lag, throughput и ошибки

Для production-pipeline необходимо отслеживать три ключевые метрики:

  • Replication lag — задержка между изменением в PostgreSQL и индексацией в Elasticsearch (в секундах).
  • Throughput — количество событий в секунду, обработанных pipeline.
  • Error rate — доля ошибок при индексации.

Экспортируем метрики через Prometheus:

package metrics

import "github.com/prometheus/client_golang/prometheus"

var (
    EventsProcessed = prometheus.NewCounterVec(
        prometheus.CounterOpts{
            Name: "pipeline_events_total",
            Help: "Total WAL events processed",
        },
        []string{"table", "operation"},
    )
    ReplicationLag = prometheus.NewGauge(
        prometheus.GaugeOpts{
            Name: "pipeline_replication_lag_seconds",
            Help: "Replication lag in seconds",
        },
    )
    IndexErrors = prometheus.NewCounter(
        prometheus.CounterOpts{
            Name: "pipeline_index_errors_total",
            Help: "Total Elasticsearch indexing errors",
        },
    )
    BatchSize = prometheus.NewHistogram(
        prometheus.HistogramOpts{
            Name:    "pipeline_batch_size",
            Help:    "Size of Elasticsearch bulk batches",
            Buckets: prometheus.LinearBuckets(10, 10, 10),
        },
    )
)

func init() {
    prometheus.MustRegister(EventsProcessed, ReplicationLag, IndexErrors, BatchSize)
}

Replication lag вычисляется как разница между текущим временем и временем коммита транзакции из WAL-события. Настройте алерты в Alertmanager: если lag превышает 30 секунд — предупреждение, 2 минуты — критическая ошибка.

Деплой в Docker и Kubernetes

Dockerfile

FROM golang:1.22-alpine AS builder
WORKDIR /app
COPY go.mod go.sum ./
RUN go mod download
COPY . .
RUN CGO_ENABLED=0 GOOS=linux go build -o pipeline ./cmd/pipeline

FROM alpine:3.19
RUN apk add --no-cache ca-certificates tzdata
WORKDIR /app
COPY --from=builder /app/pipeline .
ENTRYPOINT ["./pipeline"]

Kubernetes Deployment

Pipeline запускается как одиночный Pod (не масштабируется горизонтально — один слот репликации на одного читателя):

apiVersion: apps/v1
kind: Deployment
metadata:
  name: data-pipeline
  namespace: production
spec:
  replicas: 1
  selector:
    matchLabels:
      app: data-pipeline
  template:
    metadata:
      labels:
        app: data-pipeline
      annotations:
        prometheus.io/scrape: "true"
        prometheus.io/port: "9090"
    spec:
      containers:
      - name: pipeline
        image: your-registry/data-pipeline:latest
        env:
        - name: PG_DSN
          valueFrom:
            secretKeyRef:
              name: pipeline-secrets
              key: pg-dsn
        - name: ES_ADDR
          valueFrom:
            secretKeyRef:
              name: pipeline-secrets
              key: es-addr
        - name: SLOT_NAME
          value: "pipeline_slot"
        - name: PUBLICATION
          value: "products_pub"
        - name: BATCH_SIZE
          value: "500"
        - name: FLUSH_INTERVAL
          value: "1s"
        resources:
          requests:
            cpu: 100m
            memory: 128Mi
          limits:
            cpu: 500m
            memory: 256Mi
        livenessProbe:
          httpGet:
            path: /health
            port: 9090
          initialDelaySeconds: 10
          periodSeconds: 30

Все конфигурационные параметры передаются через переменные окружения — никаких хардкодированных значений в образе. Секреты хранятся в Kubernetes Secrets и монтируются через secretKeyRef.

Тестирование pipeline

Юнит-тесты трансформации

Функции маппинга тестируются без внешних зависимостей:

package transform_test

import (
    "testing"
    "time"

    "github.com/stretchr/testify/assert"
    "your/project/internal/transform"
)

func TestMapProductToDoc(t *testing.T) {
    row := transform.ProductRow{
        ID:          42,
        Name:        "Test Product",
        Description: "Description",
        Price:       99.99,
        Category:    "Electronics",
        Tags:        []string{"new", "sale"},
        UpdatedAt:   time.Date(2026, 1, 1, 0, 0, 0, 0, time.UTC),
    }
    doc := transform.MapProductToDoc(row)
    assert.Equal(t, "42", doc.ID)
    assert.Equal(t, "Test Product", doc.Name)
    assert.Equal(t, 99.99, doc.Price)
    assert.Equal(t, []string{"new", "sale"}, doc.Tags)
}

Интеграционные тесты

Для интеграционных тестов используйте testcontainers-go: поднимайте реальные PostgreSQL и Elasticsearch в Docker-контейнерах прямо в тестах:

func TestPipelineIntegration(t *testing.T) {
    ctx := context.Background()

    pgContainer, err := testcontainers.GenericContainer(ctx,
        testcontainers.GenericContainerRequest{
            ContainerRequest: testcontainers.ContainerRequest{
                Image:        "postgres:16",
                ExposedPorts: []string{"5432/tcp"},
                Env: map[string]string{
                    "POSTGRES_PASSWORD": "test",
                    "POSTGRES_DB":       "testdb",
                },
                Cmd: []string{"postgres", "-c", "wal_level=logical"},
                WaitingFor: wait.ForListeningPort("5432/tcp"),
            },
            Started: true,
        },
    )
    require.NoError(t, err)
    defer pgContainer.Terminate(ctx)

    // ... аналогично для Elasticsearch
    // Запустите pipeline, сделайте INSERT в PostgreSQL
    // Проверьте наличие документа в Elasticsearch через polling с таймаутом
}

Тест должен: создать таблицу и publication, запустить pipeline как горутину, выполнить INSERT/UPDATE/DELETE в PostgreSQL и через несколько секунд проверить состояние индекса Elasticsearch. Это покрывает весь критический путь.

Заключение

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

Ключевые принципы, которые делают pipeline production-ready:

  • CDC через WAL вместо polling — минимальная задержка и поддержка всех операций.
  • Сохранение LSN после успешного flush — правильная семантика at-least-once.
  • Идемпотентная индексация через явный _id — безопасность при повторах.
  • Батчевая отправка с flush по размеру и таймауту — баланс между задержкой и throughput.
  • Экспорт метрик в Prometheus — видимость состояния pipeline в реальном времени.
  • Конфигурация через переменные окружения — переносимость между окружениями.

Следующий шаг для масштабирования: если объём событий превышает возможности одного слота репликации, рассмотрите партиционирование таблиц и несколько publication с отдельными pipeline-инстансами, либо переход к Kafka как промежуточному брокеру с Debezium-коннектором.

Технологии

Теги

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

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