Construcción de un Data Pipeline robusto con Go, PostgreSQL y Elasticsearch: de los datos brutos a la búsqueda en tiempo real
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 = 4Cree 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└── DockerfileLectura 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: 30Todos 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
_idexplí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í →