Keyboard shortcuts

Press ← or → to navigate between chapters

Press S or / to search in the book

Press ? to show this help

Press Esc to hide this help

Elasticsearch

elasticsearch is an output that writes documents through the _bulk API: an action line before each document, and a response with one entry per document, so a rejected document fails alone. Tested against Elasticsearch 8.19. OpenSearch has the same _bulk API but has not been tried.

mqb copy --drain --batch-size 5000 \
  'file://books.jsonl?format=raw' \
  'elasticsearch://localhost:9200/books?api_key=<api key>'

elasticsearch+https://host/books connects over HTTPS, and ?mode=update merges each payload into the stored document instead of replacing it.

Keep an index in sync with a Postgres table

books_to_search:
  batch_size: 5000
  input:
    postgres_cdc:
      url: "postgres://user:pass@localhost/app"
      publication: "books_pub"
      slot_name: "mqb_elastic"
      consume: capture_all                 # backfill first, then stream changes
      cursor_id: "books_backfill"
      checkpoint_store: "file:///var/lib/mq-bridge/books-phase.json"
  output:
    middlewares:
      - retry: { max_attempts: 5 }
    custom:
      name: elasticsearch
      config:
        url: "elasticsearch://localhost:9200"
        api_key: "<api key>"
        index: books
        operation: "${metadata:postgres.operation}"
        compression: gzip

The document _id is the row’s id column, so an update replaces the document and a delete finds it. Name another column with id_field. See Postgres CDC for the publication and the slot.

Options

FieldDefaultMeaning
urlrequiredelasticsearch://host:9200, elasticsearch+https://host or an http(s):// URL
indexrequiredIndex to write to; the path of the URI
api_keynoneSent as Authorization: ApiKey <api key>
id_fieldidTop-level payload field that becomes the document _id
modeindexindex replaces the document; update merges the payload into it and creates it if missing
authnoneoauth2 or aws_sigv4, as on http_bulk; in a URI, JSON: auth={"aws_sigv4":{"region":"eu-central-1","service":"es"}}
operationnoneTemplate for a message’s operation; delete or d removes the document
compressionnonegzip, zstd or lz4 request bodies
request_timeout_msnoneRequest timeout

elasticsearch is the generic http_bulk output with these requests filled in. For anything it does not offer, such as basic authentication or another bulk action, write the http_bulk form:

output:
  http_bulk:
    url: http://localhost:9200
    headers:
      Authorization: ApiKey <api key>
    operation: "${metadata:postgres.operation}"
    compression: gzip
    upsert:
      path: /books/_bulk
      action: '{"index":{"_id":"${payload:id}"}}'
      result:
        items:
          path: /items
          error: /index/error/reason
    delete:
      path: /books/_bulk
      line: '{"delete":{"_id":{id}}}'
      result:
        items:
          path: /items
          error: /delete/error/reason

What to know

  • A rejected document fails alone, with Elasticsearch’s reason, for example “failed to parse field [year] of type [integer]”. Add a dlq to keep it.
  • index replaces the whole document. mode: update sends {"update":{"_id":…}} and {"doc":<payload>,"doc_as_upsert":true} instead: fields the payload leaves out keep their stored value. It was run against a stub server only, not against Elasticsearch.
  • auth was run against a stub server only, for OAuth2 and for SigV4 (Amazon OpenSearch Service).
  • The test ran with security disabled. api_key is sent in the form Elasticsearch documents for an API key.
  • Mappings and index settings are not managed. Create the index first if dynamic mapping is not what you want.
  • Another data path already exists for Kafka users: Kafka Connect’s Elasticsearch sink reads a topic that mq-bridge writes. This endpoint is for writing without Kafka in between.

The behaviour on errors is described on the http_bulk page.