Arquitectura

Elasticsearch Percolator: notificaciones inteligentes y búsqueda inversa en arquitectura de microservicios

Ruslan Ismailov Publicado 9 min de lectura
E

Introducción: qué es el Percolator y en qué se diferencia la búsqueda inversa

En el modelo de búsqueda clásico, almacenas documentos y ejecutas consultas sobre ellos. Elasticsearch Percolator invierte esta lógica: almacenas las consultas en el índice y luego compruebas cuáles de ellas coinciden con un documento entrante. Esto es lo que se denomina búsqueda inversa (reverse search).

Ejemplo práctico: un usuario configura una alerta: «notifícame cuando el iPhone 16 Pro esté disponible en stock por menos de 800 euros». Este filtro se guarda como una consulta percolate. Cuando llega un nuevo documento de producto al sistema, Elasticsearch lo comprueba contra todas las consultas almacenadas y devuelve la lista de las que han coincidido, es decir, los usuarios a los que hay que notificar.

Principales escenarios de uso del Percolator en 2026:

  • Sistemas de alertas de precios en e-commerce
  • Monitorización de flujos de noticias y agregadores RSS
  • Alertas por métricas en plataformas de observabilidad
  • Filtrado de eventos de seguridad (sistemas SIEM)
  • Suscripciones inteligentes a contenido en plataformas de medios

La diferencia clave respecto al enfoque clásico: el número de consultas almacenadas puede ser de millones, mientras que los documentos entrantes pueden ser solo unos pocos por segundo. El Percolator está optimizado precisamente para esta inversión en la relación entre consultas y datos.

Visión arquitectónica: Percolator en un sistema de microservicios

En una arquitectura de microservicios, el Percolator generalmente actúa como núcleo del servicio de notificaciones o del motor de alertas. El esquema típico de interacción incluye varias capas:

  1. Servicio de gestión de suscripciones — recibe los filtros del usuario a través de una REST API y los almacena como consultas percolate en Elasticsearch.
  2. Broker de eventos (Kafka, RabbitMQ, Pulsar) — recibe los eventos entrantes del catálogo de productos, el flujo de noticias o el sistema de monitorización.
  3. Percolate worker — microservicio consumidor que recoge los eventos del broker, ejecuta la consulta percolate contra Elasticsearch y pasa la lista de suscripciones coincidentes al servicio de notificaciones.
  4. Servicio de notificaciones — envía notificaciones por email, push, Telegram, Slack y otros canales.

El procesamiento asíncrono a través del broker de eventos es fundamental: permite percolatar documentos de forma independiente al flujo principal del producto y sin bloquear la escritura. El percolate worker puede escalar horizontalmente: cada instancia procesa su propia partición del topic.

Principio arquitectónico importante: el índice de suscripciones y el índice de datos deben estar separados. El Percolator trabaja sobre el mapping del índice de destino, pero almacena las consultas en un campo separado de tipo percolator.

Configuración del índice Percolator: mapping y almacenamiento de consultas

Veamos la configuración con el ejemplo de un índice de alertas de precios para un marketplace. Primero creamos el índice con el mapping correcto:

PUT /price-alerts
{
  "mappings": {
    "properties": {
      "query": {
        "type": "percolator"
      },
      "user_id": {
        "type": "keyword"
      },
      "notification_channel": {
        "type": "keyword"
      },
      "created_at": {
        "type": "date"
      },
      "category": {
        "type": "keyword"
      },
      "max_price": {
        "type": "double"
      },
      "product_name": {
        "type": "text",
        "analyzer": "standard"
      },
      "in_stock": {
        "type": "boolean"
      },
      "sku": {
        "type": "keyword"
      }
    }
  },
  "settings": {
    "number_of_shards": 3,
    "number_of_replicas": 1,
    "index.percolator.map_unmapped_fields_as_text": true
  }
}

El campo query con tipo percolator es el elemento clave. El resto de campos describen la estructura de los documentos que serán percolados. Esto es importante: el mapping debe corresponder a la estructura de los datos entrantes, no a las propias suscripciones.

Guardamos una alerta de usuario — una consulta para un iPhone por menos de 800 euros y en stock:

PUT /price-alerts/_doc/alert-user-42-iphone
{
  "user_id": "42",
  "notification_channel": "email",
  "created_at": "2026-01-15T10:00:00Z",
  "query": {
    "bool": {
      "must": [
        {
          "match": {
            "product_name": "iPhone 16 Pro"
          }
        },
        {
          "term": {
            "in_stock": true
          }
        }
      ],
      "filter": [
        {
          "range": {
            "max_price": {
              "lte": 80000
            }
          }
        }
      ]
    }
  }
}

Cabe destacar: dentro del campo query se puede usar cualquier consulta estándar de Elasticsearch — bool, term, match, range, geo_distance, nested y otras. Esto ofrece una enorme flexibilidad para crear filtros de usuario complejos.

Implementación del servicio de notificaciones: procesamiento de documentos

Cuando llega un producto nuevo o actualizado al catálogo, el percolate worker ejecuta la consulta de coincidencia:

POST /price-alerts/_search
{
  "query": {
    "percolate": {
      "field": "query",
      "document": {
        "product_name": "Apple iPhone 16 Pro 256GB",
        "sku": "APPL-IP16P-256",
        "max_price": 74990,
        "in_stock": true,
        "category": "smartphones"
      }
    }
  },
  "_source": ["user_id", "notification_channel"]
}

Elasticsearch devuelve todos los documentos-consulta que coinciden con el documento proporcionado. La respuesta contiene los _id de las alertas y los metadatos de los usuarios. A continuación, el worker pasa la lista al servicio de notificaciones:

# Pseudocódigo del worker en Python
def process_product_event(product: dict):
    response = es_client.search(
        index="price-alerts",
        body={
            "query": {
                "percolate": {
                    "field": "query",
                    "document": product
                }
            },
            "_source": ["user_id", "notification_channel"],
            "size": 1000  # máximo de alertas por consulta
        }
    )

    matched_alerts = response["hits"]["hits"]

    for alert in matched_alerts:
        user_id = alert["_source"]["user_id"]
        channel = alert["_source"]["notification_channel"]
        notification_queue.publish({
            "user_id": user_id,
            "channel": channel,
            "product": product,
            "alert_id": alert["_id"]
        })

    return len(matched_alerts)

Punto importante: cuando el número de coincidencias es alto, utiliza el parámetro size y, si es necesario, paginación con search_after. Por defecto, Elasticsearch solo devuelve las 10 consultas coincidentes.

Integración con microservicios a través de REST API

El servicio de gestión de suscripciones expone una REST API para el frontend y otros microservicios. Contrato típico:

# Creación de una alerta
POST /api/v1/alerts
Content-Type: application/json
Authorization: Bearer {token}

{
  "user_id": "42",
  "notification_channel": "push",
  "filters": {
    "product_name": "MacBook Pro",
    "max_price": 150000,
    "in_stock": true,
    "category": "laptops"
  }
}

# Respuesta
{
  "alert_id": "alert-user-42-macbook-001",
  "status": "active",
  "created_at": "2026-03-10T12:00:00Z"
}

El servicio de suscripciones traduce los filtros del usuario a una consulta DSL de Elasticsearch y los almacena en el índice percolator. El punto clave es validar la consulta antes de guardarla. Elasticsearch proporciona la Validate API:

POST /price-alerts/_validate/query
{
  "query": {
    "bool": {
      "must": [
        { "match": { "product_name": "MacBook Pro" } },
        { "term": { "in_stock": true } }
      ],
      "filter": [
        { "range": { "max_price": { "lte": 150000 } } }
      ]
    }
  }
}

Para el procesamiento asíncrono de eventos entrantes, usa Kafka con particionamiento por category del producto. Esto permite que cada instancia del percolate worker procese su categoría en paralelo, sin competir por las mismas alertas.

Para actualizar o eliminar una alerta basta con las operaciones estándar de Elasticsearch Update/Delete: el índice percolator se comporta como un índice normal desde el punto de vista del CRUD.

Rendimiento: características de carga y optimización

El rendimiento del Percolator está determinado principalmente por el número de consultas almacenadas y su complejidad. Métricas y recomendaciones clave:

  • Número de shards: Elasticsearch ejecuta el percolate en paralelo en todos los shards. Con un millón de consultas, lo óptimo es entre 5 y 10 shards. Un sharding excesivo genera overhead de coordinación.
  • Caché: el Percolator hace un uso intensivo del query cache. Los términos y filtros (term, range) se cachean eficientemente. Las consultas match se cachean peor — utiliza term siempre que sea posible.
  • Ordenación por score: si no necesitas el relevance score, añade "sort": ["_doc"] — esto acelera la recuperación entre un 20 y un 40%.
  • Tamaño del documento: cuanto más pequeño sea el documento que percolatas, más rápido será el procesamiento. No transmitas campos innecesarios.
  • Named queries: usa _name en las consultas para depuración, pero desactívalo en producción, ya que añade overhead.

Ejemplo de consulta con optimizaciones para producción:

POST /price-alerts/_search
{
  "query": {
    "percolate": {
      "field": "query",
      "document": {
        "product_name": "Samsung Galaxy S25",
        "max_price": 65000,
        "in_stock": true,
        "category": "smartphones"
      }
    }
  },
  "sort": ["_doc"],
  "_source": ["user_id", "notification_channel"],
  "size": 500,
  "track_total_hits": false
}

Los benchmarks en un clúster de 3 nodos con 16 CPU / 64 GB RAM muestran: con 500 000 consultas percolate almacenadas, el tiempo medio de ejecución de un percolate es de 15 a 50 ms según la complejidad de las consultas. Con 5 millones de consultas, el tiempo es de 100 a 300 ms, lo que requiere una cuidadosa optimización del sharding y del hardware.

Caso de uso real: sistema de alertas de precios end-to-end

Veamos el ciclo completo de funcionamiento de un sistema de monitorización de precios para un marketplace con 2 millones de usuarios:

  1. El usuario crea una alerta a través de la aplicación móvil: «Notifícame sobre el MacBook Pro por menos de 1500 euros».
  2. El Subscription Service valida la consulta, la traduce al DSL de Elasticsearch y la guarda en el índice price-alerts. Simultáneamente, registra la metainformación (email, push token) en PostgreSQL.
  3. El proveedor actualiza el precio a través del Catalog Service → el evento se publica en el topic de Kafka product.price.updated.
  4. El Percolate Worker (microservicio en Go) consume el evento, ejecuta la consulta percolate y obtiene la lista de alertas coincidentes.
  5. Deduplicación: los resultados se verifican a través de Redis — no notificar al mismo usuario sobre la misma alerta más de una vez cada 24 horas. Clave: notif:{alert_id}:{date}.
  6. El Notification Worker envía notificaciones push a través de Firebase y emails a través de SendGrid. El estado se registra en PostgreSQL.
// Go: Percolate Worker — ejemplo simplificado
func (w *PercolateWorker) HandleProductEvent(ctx context.Context, product Product) error {
    res, err := w.esClient.Search(
        w.esClient.Search.WithIndex("price-alerts"),
        w.esClient.Search.WithBody(strings.NewReader(fmt.Sprintf(`{
            "query": {
                "percolate": {
                    "field": "query",
                    "document": {
                        "product_name": %q,
                        "max_price": %f,
                        "in_stock": %v,
                        "category": %q
                    }
                }
            },
            "sort": ["_doc"],
            "_source": ["user_id", "notification_channel"],
            "size": 1000,
            "track_total_hits": false
        }`, product.Name, product.Price, product.InStock, product.Category))),
    )
    if err != nil {
        return fmt.Errorf("percolate query failed: %w", err)
    }
    defer res.Body.Close()

    var result PercolateResult
    if err := json.NewDecoder(res.Body).Decode(&result); err != nil {
        return fmt.Errorf("decode response failed: %w", err)
    }

    for _, hit := range result.Hits.Hits {
        dedupKey := fmt.Sprintf("notif:%s:%s", hit.ID, time.Now().Format("2006-01-02"))
        if w.redis.SetNX(ctx, dedupKey, 1, 24*time.Hour).Val() {
            w.notifQueue.Publish(ctx, NotificationTask{
                AlertID:  hit.ID,
                UserID:   hit.Source.UserID,
                Channel:  hit.Source.NotificationChannel,
                Product:  product,
            })
        }
    }
    return nil
}

Esta arquitectura procesa hasta 10 000 eventos de precios por minuto con 2 millones de alertas activas, cumpliendo un SLA de 500 ms end-to-end desde el evento hasta el envío de la notificación.

Problemas y limitaciones del Percolator

El Percolator es una herramienta potente, pero tiene limitaciones importantes que deben tenerse en cuenta en el diseño:

  • No soporta consultas join entre documentos. El Percolator comprueba un documento a la vez. Si tu alerta debe considerar datos relacionados de otro índice, tendrás que desnormalizar los datos antes de percolatar.
  • La complejidad de las consultas tiene un impacto no lineal. Las consultas bool anidadas con wildcard o filtros de script pueden ralentizar el procesamiento en varios órdenes de magnitud. Perfila con "profile": true.
  • La actualización del mapping requiere reindexación. Si la estructura de los documentos cambia, hay que recrear el índice de alertas y redirigir las consultas almacenadas.
  • No hay deduplicación nativa. Elasticsearch no rastrea qué alertas ya han sido disparadas — esa es tu responsabilidad (Redis, PostgreSQL).
  • No se admiten consultas con agregaciones. Una consulta percolate no puede contener agregaciones dentro del filtro almacenado.
  • No es adecuado para eventos de alta frecuencia con reglas simples. Si tienes 100 000 eventos por segundo con condiciones simples, considera soluciones como Kafka Streams o Flink, que ofrecen menor latencia.
  • Tamaño del índice. Millones de consultas complejas requieren una cantidad significativa de RAM para los segmentos de Lucene. Planifica el heap de Elasticsearch a razón de ~1–5 KB por consulta.

El Percolator es ideal cuando el número de filtros únicos de usuario supera en órdenes de magnitud la frecuencia de los eventos entrantes, y los propios filtros tienen una complejidad moderada. Para reglas simples con alta frecuencia de eventos, considera el procesamiento de streams; para condiciones join muy complejas, un rule engine a nivel de aplicación.

Conclusión

Elasticsearch Percolator sigue siendo una de las herramientas más elegantes para construir sistemas de notificaciones inteligentes y búsqueda inversa en arquitecturas de microservicios. En 2026 es más relevante que nunca: el crecimiento del número de filtros personalizados en e-commerce, medios y observabilidad genera exactamente el patrón de carga para el que el Percolator está optimizado.

Conclusiones clave para la aplicación práctica:

  • Separa los índices de datos y las consultas percolator, y usa el mapping correcto desde el principio.
  • Construye el procesamiento asíncrono a través de un broker de eventos — es un requisito imprescindible para un sistema listo para producción.
  • Implementa la deduplicación de notificaciones a nivel de Redis, fuera de Elasticsearch.
  • Perfila las consultas complejas y da preferencia a los filtros term/range frente a match y script.
  • Prueba el rendimiento con el volumen real de consultas almacenadas — la degradación es no lineal.

Un sistema bien diseñado con Elasticsearch Percolator puede gestionar millones de alertas de usuario en tiempo real con latencias aceptables y escalado horizontal — sin necesidad de escribir tu propio rule engine desde cero.

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í →