Построение надёжного Data Pipeline с Go, PostgreSQL и Elasticsearch: от сырых данных до поиска в реальном времени
Введение: задача синхронизации данных
Современные приложения всё чаще комбинируют реляционные базы данных с полнотекстовыми поисковыми движками. 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. Подробнее обо мне →