Desarrollo backend

Construcción de un Data Pipeline robusto con Go, PostgreSQL y Elasticsearch: de los datos brutos a la búsqueda en tiempo real

Ruslan Ismailov Publicado 14 min de lectura
C

Introducción: el desafío de la sincronización de datos

Las aplicaciones modernas combinan cada vez más bases de datos relacionales con motores de búsqueda de texto completo. PostgreSQL almacena datos transaccionales de forma fiable, pero sus capacidades de búsqueda de texto completo son limitadas en comparación con Elasticsearch. La sincronización entre PostgreSQL y Elasticsearch es uno de los principales retos arquitectónicos para los desarrolladores backend en 2026.

El enfoque directo de "escribir simultáneamente en ambos sistemas desde la aplicación" no es fiable: si una operación falla, los datos divergen. La solución correcta es un data pipeline dedicado que rastree los cambios en PostgreSQL y los replique de forma atómica en Elasticsearch. Precisamente ese pipeline construiremos en este artículo usando Go.

Arquitectura del pipeline: polling, triggers y CDC

Existen tres enfoques principales para sincronizar datos entre PostgreSQL y Elasticsearch:

  • Polling — consulta periódica de tablas mediante el campo updated_at. Sencillo de implementar, pero introduce latencia y no detecta eliminaciones.
  • Database Triggers — los triggers de PostgreSQL registran los cambios en una tabla outbox. Más fiable que el polling, pero incrementa la carga sobre la base de datos.
  • Change Data Capture (CDC) — lectura del journal WAL (Write-Ahead Log) de PostgreSQL. Carga mínima, soporte de todas las operaciones (INSERT, UPDATE, DELETE), latencia en milisegundos.

Para sistemas en producción, la opción óptima es CDC mediante replicación lógica de PostgreSQL. Es el único enfoque que garantiza la integridad de los datos sin carga adicional sobre la aplicación. Es el que implementaremos con Go.

Change Data Capture con PostgreSQL: configuración de la replicación lógica

Para trabajar con CDC es necesario habilitar la replicación lógica en PostgreSQL. Edite postgresql.conf:

wal_level = logical\nmax_replication_slots = 4\nmax_wal_senders = 4

Cree un slot de replicación y una publication para las tablas necesarias:

-- Creación de la publication\nCREATE PUBLICATION products_pub FOR TABLE products, categories;\n\n-- Creación del slot de replicación lógica\nSELECT pg_create_logical_replication_slot('pipeline_slot', 'pgoutput');

Para leer eventos WAL en Go utilizamos la biblioteca pglogrepl — un cliente Go puro para el protocolo de replicación lógica de PostgreSQL. La alternativa es Debezium (solución JVM), conveniente si la infraestructura ya usa Kafka. En nuestro caso, Go + pglogrepl ofrece dependencias mínimas y máximo control.

Implementación del pipeline en Go

Estructura del proyecto

Organizamos el código siguiendo el principio de arquitectura limpia:

pipeline/\n├── cmd/\n│   └── pipeline/\n│       └── main.go\n├── internal/\n│   ├── cdc/\n│   │   ├── reader.go       # Lectura de eventos WAL\n│   │   └── decoder.go      # Decodificación de pgoutput\n│   ├── transform/\n│   │   └── mapper.go       # Mapeo PG → documentos ES\n│   ├── indexer/\n│   │   └── elasticsearch.go # Indexación por lotes\n│   └── metrics/\n│       └── prometheus.go   # Métricas\n├── docker-compose.yml\n└── Dockerfile

Lectura de eventos WAL

El worker principal se conecta a PostgreSQL mediante el protocolo de replicación y lee el flujo de cambios:

package cdc\n\nimport (\n    "context"\n    "fmt"\n    "time"\n\n    "github.com/jackc/pglogrepl"\n    "github.com/jackc/pgx/v5/pgconn"\n    "github.com/jackc/pgx/v5/pgproto3"\n)\n\ntype WALReader struct {\n    conn        *pgconn.PgConn\n    slotName    string\n    publication string\n    lsn         pglogrepl.LSN\n}\n\nfunc NewWALReader(dsn, slotName, publication string) (*WALReader, error) {\n    conn, err := pgconn.Connect(context.Background(), dsn+" replication=database")\n    if err != nil {\n        return nil, fmt.Errorf("connect replication: %w", err)\n    }\n    return &WALReader{\n        conn:        conn,\n        slotName:    slotName,\n        publication: publication,\n    }, nil\n}\n\nfunc (r *WALReader) Start(ctx context.Context, events chan<- *WALEvent) error {\n    opts := pglogrepl.StartReplicationOptions{\n        PluginArgs: []string{\n            "proto_version '1'",\n            fmt.Sprintf("publication_names '%s'", r.publication),\n        },\n    }\n    if err := pglogrepl.StartReplication(ctx, r.conn, r.slotName, r.lsn, opts); err != nil {\n        return fmt.Errorf("start replication: %w", err)\n    }\n\n    standbyDeadline := time.Now().Add(10 * time.Second)\n    for {\n        if time.Now().After(standbyDeadline) {\n            if err := pglogrepl.SendStandbyStatusUpdate(ctx, r.conn,\n                pglogrepl.StandbyStatusUpdate{WALWritePosition: r.lsn}); err != nil {\n                return fmt.Errorf("standby status: %w", err)\n            }\n            standbyDeadline = time.Now().Add(10 * time.Second)\n        }\n\n        ctx2, cancel := context.WithDeadline(ctx, standbyDeadline)\n        msg, err := r.conn.ReceiveMessage(ctx2)\n        cancel()\n        if err != nil {\n            if pgconn.Timeout(err) {\n                continue\n            }\n            return fmt.Errorf("receive message: %w", err)\n        }\n\n        switch m := msg.(type) {\n        case *pgproto3.CopyData:\n            if m.Data[0] == pglogrepl.XLogDataByteID {\n                xld, err := pglogrepl.ParseXLogData(m.Data[1:])\n                if err != nil {\n                    return fmt.Errorf("parse xlog: %w", err)\n                }\n                event, err := DecodeWALData(xld.WALData)\n                if err == nil && event != nil {\n                    events <- event\n                    r.lsn = xld.WALStart + pglogrepl.LSN(len(xld.WALData))\n                }\n            }\n        }\n    }\n}

Indexación por lotes en Elasticsearch

Para una indexación eficiente usamos la Bulk API de Elasticsearch. Acumulamos eventos en un buffer y los enviamos según el tamaño o un temporizador:

package indexer\n\nimport (\n    "bytes"\n    "context"\n    "encoding/json"\n    "fmt"\n    "time"\n\n    "github.com/elastic/go-elasticsearch/v8"\n    "github.com/elastic/go-elasticsearch/v8/esapi"\n)\n\ntype BulkIndexer struct {\n    client    *elasticsearch.Client\n    index     string\n    batchSize int\n    flushInterval time.Duration\n    buf       []BulkAction\n}\n\ntype BulkAction struct {\n    ID     string\n    Doc    map[string]interface{}\n    Delete bool\n}\n\nfunc (bi *BulkIndexer) Run(ctx context.Context, actions <-chan BulkAction) error {\n    ticker := time.NewTicker(bi.flushInterval)\n    defer ticker.Stop()\n\n    for {\n        select {\n        case action, ok := <-actions:\n            if !ok {\n                return bi.flush(ctx)\n            }\n            bi.buf = append(bi.buf, action)\n            if len(bi.buf) >= bi.batchSize {\n                if err := bi.flush(ctx); err != nil {\n                    return err\n                }\n            }\n        case <-ticker.C:\n            if len(bi.buf) > 0 {\n                if err := bi.flush(ctx); err != nil {\n                    return err\n                }\n            }\n        case <-ctx.Done():\n            return ctx.Err()\n        }\n    }\n}\n\nfunc (bi *BulkIndexer) flush(ctx context.Context) error {\n    if len(bi.buf) == 0 {\n        return nil\n    }\n    var body bytes.Buffer\n    for _, action := range bi.buf {\n        if action.Delete {\n            meta := map[string]interface{}{"delete": map[string]interface{}{"_index": bi.index, "_id": action.ID}}\n            line, _ := json.Marshal(meta)\n            body.Write(line)\n            body.WriteByte('\n')\n        } else {\n            meta := map[string]interface{}{"index": map[string]interface{}{"_index": bi.index, "_id": action.ID}}\n            line, _ := json.Marshal(meta)\n            body.Write(line)\n            body.WriteByte('\n')\n            doc, _ := json.Marshal(action.Doc)\n            body.Write(doc)\n            body.WriteByte('\n')\n        }\n    }\n    req := esapi.BulkRequest{Body: &body}\n    res, err := req.Do(ctx, bi.client)\n    if err != nil {\n        return fmt.Errorf("bulk request: %w", err)\n    }\n    defer res.Body.Close()\n    if res.IsError() {\n        return fmt.Errorf("bulk response error: %s", res.Status())\n    }\n    bi.buf = bi.buf[:0]\n    return nil\n}

Garantías de entrega: at-least-once e idempotencia

La replicación WAL en PostgreSQL garantiza la entrega at-least-once: al reiniciar el pipeline, los eventos pueden leerse de nuevo. La Bulk API de Elasticsearch con un _id explícito garantiza idempotencia — reindexar el mismo documento es seguro.

Es fundamental guardar el LSN (Log Sequence Number) tras el envío exitoso del lote a Elasticsearch, no tras la lectura del WAL. Almacene el LSN en una tabla separada de PostgreSQL o en Redis:

// Guardado del LSN tras un flush exitoso\nfunc saveLSN(ctx context.Context, db *pgxpool.Pool, slotName string, lsn pglogrepl.LSN) error {\n    _, err := db.Exec(ctx,\n        `INSERT INTO pipeline_checkpoints (slot_name, lsn, updated_at)\n         VALUES ($1, $2, now())\n         ON CONFLICT (slot_name) DO UPDATE SET lsn = $2, updated_at = now()`,\n        slotName, lsn.String(),\n    )\n    return err\n}

El manejo de errores durante el flush debe incluir backoff exponencial con jitter para no sobrecargar Elasticsearch durante fallos temporales:

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

Transformación de datos: mapeo PostgreSQL → Elasticsearch

Las estructuras de datos en PostgreSQL y los documentos de Elasticsearch suelen diferir: las tablas normalizadas deben desnormalizarse, los tipos deben convertirse y los campos deben filtrarse o enriquecerse.

package transform\n\nimport "time"\n\n// Fila de PostgreSQL\ntype ProductRow struct {\n    ID          int64\n    Name        string\n    Description string\n    Price       float64\n    CategoryID  int64\n    Category    string\n    Tags        []string\n    CreatedAt   time.Time\n    UpdatedAt   time.Time\n    Deleted     bool\n}\n\n// Documento de Elasticsearch\ntype ProductDoc struct {\n    ID          string    `json:"id"`\n    Name        string    `json:"name"`\n    Description string    `json:"description"`\n    Price       float64   `json:"price"`\n    Category    string    `json:"category"`\n    Tags        []string  `json:"tags"`\n    UpdatedAt   time.Time `json:"updated_at"`\n}\n\nfunc MapProductToDoc(row ProductRow) ProductDoc {\n    return ProductDoc{\n        ID:          fmt.Sprintf("%d", row.ID),\n        Name:        row.Name,\n        Description: row.Description,\n        Price:       row.Price,\n        Category:    row.Category,\n        Tags:        row.Tags,\n        UpdatedAt:   row.UpdatedAt,\n    }\n}\n\nfunc DocToMap(doc ProductDoc) map[string]interface{} {\n    return map[string]interface{}{\n        "id":          doc.ID,\n        "name":        doc.Name,\n        "description": doc.Description,\n        "price":       doc.Price,\n        "category":    doc.Category,\n        "tags":        doc.Tags,\n        "updated_at":  doc.UpdatedAt,\n    }\n}

Tenga en cuenta que al recibir un evento DELETE desde el WAL, el pipeline debe eliminar el documento en Elasticsearch por su _id, sin consultar PostgreSQL (la fila ya fue eliminada).

Monitorización del pipeline: métricas de lag, throughput y errores

Para un pipeline en producción es necesario rastrear tres métricas clave:

  • Replication lag — retraso entre el cambio en PostgreSQL y la indexación en Elasticsearch (en segundos).
  • Throughput — número de eventos por segundo procesados por el pipeline.
  • Error rate — proporción de errores durante la indexación.

Exportamos las métricas a través de Prometheus:

package metrics\n\nimport "github.com/prometheus/client_golang/prometheus"\n\nvar (\n    EventsProcessed = prometheus.NewCounterVec(\n        prometheus.CounterOpts{\n            Name: "pipeline_events_total",\n            Help: "Total WAL events processed",\n        },\n        []string{"table", "operation"},\n    )\n    ReplicationLag = prometheus.NewGauge(\n        prometheus.GaugeOpts{\n            Name: "pipeline_replication_lag_seconds",\n            Help: "Replication lag in seconds",\n        },\n    )\n    IndexErrors = prometheus.NewCounter(\n        prometheus.CounterOpts{\n            Name: "pipeline_index_errors_total",\n            Help: "Total Elasticsearch indexing errors",\n        },\n    )\n    BatchSize = prometheus.NewHistogram(\n        prometheus.HistogramOpts{\n            Name:    "pipeline_batch_size",\n            Help:    "Size of Elasticsearch bulk batches",\n            Buckets: prometheus.LinearBuckets(10, 10, 10),\n        },\n    )\n)\n\nfunc init() {\n    prometheus.MustRegister(EventsProcessed, ReplicationLag, IndexErrors, BatchSize)\n}

El replication lag se calcula como la diferencia entre la hora actual y la hora de commit de la transacción del evento WAL. Configure alertas en Alertmanager: si el lag supera los 30 segundos — advertencia; si supera los 2 minutos — error crítico.

Despliegue en Docker y Kubernetes

Dockerfile

FROM golang:1.22-alpine AS builder\nWORKDIR /app\nCOPY go.mod go.sum ./\nRUN go mod download\nCOPY . .\nRUN CGO_ENABLED=0 GOOS=linux go build -o pipeline ./cmd/pipeline\n\nFROM alpine:3.19\nRUN apk add --no-cache ca-certificates tzdata\nWORKDIR /app\nCOPY --from=builder /app/pipeline .\nENTRYPOINT ["./pipeline"]

Kubernetes Deployment

El pipeline se ejecuta como un Pod único (no escala horizontalmente — un slot de replicación por un único lector):

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

Todos los parámetros de configuración se pasan mediante variables de entorno — sin valores hardcodeados en la imagen. Los secretos se almacenan en Kubernetes Secrets y se montan mediante secretKeyRef.

Pruebas del pipeline

Tests unitarios de la transformación

Las funciones de mapeo se prueban sin dependencias externas:

package transform_test\n\nimport (\n    "testing"\n    "time"\n\n    "github.com/stretchr/testify/assert"\n    "your/project/internal/transform"\n)\n\nfunc TestMapProductToDoc(t *testing.T) {\n    row := transform.ProductRow{\n        ID:          42,\n        Name:        "Test Product",\n        Description: "Description",\n        Price:       99.99,\n        Category:    "Electronics",\n        Tags:        []string{"new", "sale"},\n        UpdatedAt:   time.Date(2026, 1, 1, 0, 0, 0, 0, time.UTC),\n    }\n    doc := transform.MapProductToDoc(row)\n    assert.Equal(t, "42", doc.ID)\n    assert.Equal(t, "Test Product", doc.Name)\n    assert.Equal(t, 99.99, doc.Price)\n    assert.Equal(t, []string{"new", "sale"}, doc.Tags)\n}

Tests de integración

Para las pruebas de integración utilice testcontainers-go: levante instancias reales de PostgreSQL y Elasticsearch en contenedores Docker directamente desde los tests:

func TestPipelineIntegration(t *testing.T) {\n    ctx := context.Background()\n\n    pgContainer, err := testcontainers.GenericContainer(ctx,\n        testcontainers.GenericContainerRequest{\n            ContainerRequest: testcontainers.ContainerRequest{\n                Image:        "postgres:16",\n                ExposedPorts: []string{"5432/tcp"},\n                Env: map[string]string{\n                    "POSTGRES_PASSWORD": "test",\n                    "POSTGRES_DB":       "testdb",\n                },\n                Cmd: []string{"postgres", "-c", "wal_level=logical"},\n                WaitingFor: wait.ForListeningPort("5432/tcp"),\n            },\n            Started: true,\n        },\n    )\n    require.NoError(t, err)\n    defer pgContainer.Terminate(ctx)\n\n    // ... de forma similar para Elasticsearch\n    // Inicie el pipeline, haga un INSERT en PostgreSQL\n    // Verifique la presencia del documento en Elasticsearch mediante polling con timeout\n}

El test debe: crear la tabla y la publication, iniciar el pipeline como goroutine, ejecutar INSERT/UPDATE/DELETE en PostgreSQL y verificar el estado del índice de Elasticsearch tras unos segundos. Esto cubre toda la ruta crítica.

Conclusión

Construir un data pipeline robusto con Go, PostgreSQL y Elasticsearch es una tarea que se resuelve con la combinación correcta de herramientas probadas: la replicación lógica de PostgreSQL garantiza la integridad de los datos, Go aporta eficiencia y simplicidad de código, y la Bulk API de Elasticsearch hace que la indexación sea escalable.

Principios clave que hacen que el pipeline sea production-ready:

  • CDC mediante WAL en lugar de polling — latencia mínima y soporte de todas las operaciones.
  • Guardar el LSN tras un flush exitoso — semántica correcta de at-least-once.
  • Indexación idempotente mediante _id explícito — seguridad ante reenvíos.
  • Envío por lotes con flush por tamaño y temporizador — equilibrio entre latencia y throughput.
  • Exportación de métricas a Prometheus — visibilidad del estado del pipeline en tiempo real.
  • Configuración mediante variables de entorno — portabilidad entre entornos.

El siguiente paso para escalar: si el volumen de eventos supera las capacidades de un único slot de replicación, considere el particionamiento de tablas con múltiples publications e instancias de pipeline independientes, o la migración a Kafka como broker intermediario con un conector Debezium.

Tecnologías

Etiquetas

Ruslan Ismailov

Desarrollador Senior Web / Backend. Desarrollador senior web/backend con 9 años de experiencia. Stack: PHP, Laravel, PostgreSQL, Redis, Docker, Kubernetes, REST, microservicios, CI/CD. Más sobre mí →