Elasticsearch Service Sink Connector for Confluent Platform

Note

  • The Elasticsearch Sink connector for Confluent Platform provides support for Elasticsearch version 8.x and later, including version 9.x. Starting with connector version 16.0.0, the connector uses the Elasticsearch Java API Client, the client library that Elastic maintains for Elasticsearch 8.x and 9.x. When the connector writes to Elasticsearch 9.x, it uses Elastic’s REST API compatibility automatically, so you don’t have to enable anything in the connector configuration.

  • Elastic’s client compatibility guarantee covers Elasticsearch versions 8.19.x and later, so the recommended versions are Elasticsearch 8.19.x or 9.x. Elasticsearch versions 8.0 through 8.18 work with the connector but they are not covered under the guarantee. Elasticsearch versions 7.x and earlier need the connector version 15.x. If you are upgrading from connector version 15.x, see Upgrade from version 15.x.

  • Effective July 6, 2025, only self-managed connector versions that meet or exceed the minimum version listed on the Supported Connector Versions page receive support from Confluent. Older, unsupported connector versions have been removed from Confluent Marketplace and are no longer available for download.

  • If you are running connector version 14.0.x, 14.1.x, 15.0.x, or 15.1.x, upgrade to this connector version, which resolves third-party Common Vulnerabilities and Exposures (CVEs) present in those versions.

The Kafka Connect Elasticsearch Service Sink connector moves data from Apache Kafka® to Elasticsearch. It writes data from a topic in Kafka to an index —a collection of documents—in Elasticsearch. All data for a topic is written to one index with one mapping. This allows an independent evolution of schemas for data from different topics. This simplifies the schema evolution because Elasticsearch has one enforcement on mappings. It means that all fields with the same name in the same index must have the same mapping.

Elasticsearch is often used for text queries, analytics and as a key-value store (use cases). The connector covers both the analytics and key-value store use cases. For the analytics use case, each message in Kafka is treated as an event and the connector uses topic+partition+offset as a unique identifier for events, which are then converted to unique documents in Elasticsearch.

For the key-value store use case, it supports using keys from Kafka messages as document IDs in Elasticsearch and provides configurations ensuring that updates to a key are written to Elasticsearch in order. For the ordering guarantee, see max.in.flight.requests. For both use cases, the connector uses Elasticsearch’s idempotent write semantics so that redelivered records don’t create duplicate documents. For details, see Delivery guarantees.

Mapping is the process of defining how a document and the fields it contains are stored and indexed. Users can explicitly define mappings for indices. When mapping is not explicitly defined, Elasticsearch can determine field names and types from data. However, types such as timestamp and decimal may not be correctly inferred. To ensure that these types are correctly inferred, the connector provides a feature to infer mapping from the schemas of Kafka messages.

Note

The Kafka source topic name is used to create the destination index name in Elasticsearch. You can change this name prior to it being used as the index name with a Single Message Transformation (SMT)–RegexRouter or TimeStampRouter–only when the flush.synchronously configuration property is set to true. For more details, see Limitations.

Features

Delivery guarantees

The connector delivers records from Kafka at least once. After a task restart or rebalance, Connect can redeliver records that were already written. The connector relies on Elasticsearch’s idempotent write semantics so that a redelivered record does not create a duplicate document. Each record is written with a deterministic document ID, that is either the record key or topic+partition+offset when key.ignore is set to true.

Whether a redelivered or out-of-order record can overwrite a newer document depends on the write path:

  • Inserts and deletes with key.ignore=false on a regular index: the connector sets the document’s external version from the Kafka record offset, or from the header named by external.version.header. Elasticsearch rejects any write whose version is not newer than the stored document, so each key converges on its latest record. This gives exactly-once, in-order results per key, at any max.in.flight.requests value.

  • Writes with key.ignore=true on a regular index: the document ID is unique per record, so a redelivery overwrites the same document. No duplicates are created.

  • Writes to data streams: the connector uses a create operation, so a record whose document ID already exists in the data stream is dropped as a duplicate and the resulting version conflict is not sent to the dead letter queue (DLQ). With key.ignore=true this affects only redeliveries. No duplicates are created.

  • Upserts (write.method=upsert): the update carries no version. A redelivered older record re-applies its fields over a newer document. Ordering between requests is guaranteed only with max.in.flight.requests=1.

  • use.autogenerated.ids=true: Elasticsearch assigns the ID, so a redelivered record creates a duplicate document. Delivery is at least once.

Dead Letter Queue

This connector supports the Dead Letter Queue (DLQ) functionality. For information about accessing and using the DLQ, see Confluent Platform Dead Letter Queue.

The following configuration properties determine whether a record is sent to the DLQ and whether the connector’s offset advances past it:

  • key.ignore: Controls whether the connector uses external versioning based on the Kafka offset when writing documents. This affects whether a version_conflict_engine_exception from Elasticsearch is treated as an expected, benign conflict that isn’t sent to the DLQ, or as an unexpected one that is sent to the DLQ).

  • behavior.on.null.values: Controls what happens to records with a null value, such as tombstones, including whether they’re reported to the DLQ and whether the offset advances.

  • behavior.on.malformed.documents: Controls what happens when Elasticsearch rejects a document, including whether the connector task fails or continues after the record is reported to the DLQ.

  • drop.invalid.message: Controls whether the connector task fails or continues after a record that failed conversion is reported to the DLQ.

For the specific behavior of each setting and value, see the Data Conversion configuration properties.

Multiple tasks

The Elasticsearch Service Sink connector supports running one or more tasks. You can specify the number of tasks in the tasks.max configuration parameter. Running more tasks increases throughput when the source topics have multiple partitions.

Mapping inference

The connector can infer mappings from Connect schemas. When enabled, the connector creates mappings based on schemas of Kafka messages. If a field is missing, the inference is limited to field types and default values. You should manually create mappings if more customizations are needed (for example, user-defined analyzers).

Schema evolution

The connector does not evolve Elasticsearch mappings itself. When schema.ignore is false, the connector infers a mapping from the schema of the first record written to an index, data stream, or alias. It creates that mapping only if the resource has none, and does not update it afterward. Fields that appear in later records are typed by Elasticsearch dynamic mapping. As a result, backward, forward, and fully compatible schema changes in Connect are handled by Elasticsearch. Some incompatible changes, such as changing a field from an integer to a string, are also handled by Elasticsearch rather than by the connector. For details, see Schema evolution.

Client-side encryption

This connector supports Client-Side Field Level Encryption (CSFLE) and Client-Side Payload Encryption (CSPE). For more information, see Manage Client-Side Encryption.

Aliases

An alias is an alternate name that points to one or more Elasticsearch indices or data streams; a data stream is a named resource that manages a sequence of backing indices for append-only, time-series data. The connector writes to aliases for both indices and data streams, but these aliases must be pre-created in Elasticsearch.

External topic to resource mapping

The connector supports external topic to resource mapping, allowing to map Kafka topics to user-defined Elasticsearch resources and write to pre-created indices, data streams, and aliases. This helps with custom naming schemes and integrating with existing Elasticsearch resources. All resources referenced via topic.to.external.resource.mapping (whether index, data stream, or alias) must exist before the connector starts. Each Kafka topic must map to only one Elasticsearch resource; many-to-one or one-to-many mappings aren’t supported.

Note

The connector supports external topic-to-resource mapping starting in version 14.1.5 and later.

Understanding external.resource.usage configuration

The external.resource.usage property dictates how the Elasticsearch connector interacts with Elasticsearch resources (indices, data streams, or aliases). Its behavior changes based on its value and the presence of other data stream-related configurations. Consider the following scenarios:

  • When external.resource.usage = DISABLED (default) and data stream configurations not set

    The connector writes to a regular Elasticsearch index, which it automatically creates using the Kafka topic name. This is the default behavior when external.resource.usage is disabled and no data stream-specific configurations are provided.

  • When external.resource.usage = DISABLED and data stream configurations provided

    If external.resource.usage is DISABLED but data stream configurations are provided, the connector automatically creates a data stream named as {type}-{dataset}-{namespace} and writes to it. This occurs when:

    • data.stream.type is not set to none.

    • data.stream.dataset is not set to none.

    • data.stream.namespace (optional) defaults to ${topic} name if not explicitly set.

    • data.stream.timestamp.field (optional) defaults to the Kafka record timestamp if not set.

    The timestamp.field is used as the @timestamp for indexing; if not set, the Kafka record timestamp is used.

  • When external.resource.usage = INDEX or ALIAS_INDEX

    Users must pre-create the target Elasticsearch index or alias-to-index. A one-to-one mapping between Kafka topics and these pre-existing resources must be provided via the topic.to.external.resource.mapping configuration (for example, payments:index-payments, logs:alias-logs). Records from each topic are then written directly to its specified index or index alias.

  • When external.resource.usage = DATASTREAM or ALIAS_DATASTREAM

    Users must pre-create the target Elasticsearch data stream or alias-to-data stream. A one-to-one topic-to-resource mapping must be defined via topic.to.external.resource.mapping (for example, metrics:metrics-ds, orders:alias-orders-ds). A common timestamp field must be configured using data.stream.timestamp.field (or the Kafka timestamp will be used by default), as all data streams require an @timestamp field.

License

The following are required to run the Kafka Connect Elasticsearch Sink connector:

The source code is available at https://github.com/confluentinc/kafka-connect-elasticsearch.

Linux on IBM Z (s390x) support

Starting with Confluent Platform 8.2, this connector supports Linux on IBM Z (s390x). The connector supports the same capability available on x86_64 unless otherwise noted. For more information, see Linux on IBM Z (s390x) support.

Configuration Properties

For a complete list of configuration properties for this connector, see Configuration Reference for Elasticsearch Service Sink Connector for Confluent Platform.

For an example of how to get Kafka Connect connected to Confluent Cloud, see Connect Self-Managed Kafka Connect to Confluent Cloud.

Limitations

  • The Elasticsearch Service Sink connector supports only Elasticsearch. It doesn’t support Amazon Elasticsearch Service, Amazon OpenSearch Service, OpenSearch, or any other backend. The connector verifies the server by checking for the X-Elastic-Product: Elasticsearch response header. If the header is missing, for example, because the server is not Elasticsearch or because a proxy or load balancer removes it, the connector configuration validation fails. A task that is already running fails on its first write after its retries are exhausted.

  • The connector only supports the following data stream types: logs and metrics. For more details, see the data.stream.type configuration property.

  • The connector does not currently support Single Message Transformations (SMTs) that modify the topic name. Additionally, the following transformations are not allowed:

    • io.debezium.transforms.ByLogicalTableRouter

    • io.debezium.transforms.outbox.EventRouter

    • org.apache.kafka.connect.transforms.RegexRouter

    • org.apache.kafka.connect.transforms.TimestampRouter

    • io.confluent.connect.transforms.MessageTimestampRouter

    • io.confluent.connect.transforms.ExtractTopic$Key

    • io.confluent.connect.transforms.ExtractTopic$Value

    Note

    These SMT limitations are inapplicable to the Elasticsearch Sink connector when the flush.synchronously configuration property is set to true. For more information about the flush.synchronously configuration property see the Configuration Reference for Elasticsearch Service Sink Connector for Confluent Platform documentation page.

  • Topic-mutating SMTs (for example, RegexRouter) aren’t supported when external.resource.usage is set to INDEX, ALIAS_INDEX, DATASTREAM, or ALIAS_DATASTREAM. These SMTs change the topic name after the connector has mapped it to a specific Elasticsearch resource, causing a mismatch with the topic.to.external.resource.mapping configuration and resulting in connector failure. For more information, see External topic to resource mapping.

Upgrade from version 15.x

Connector version 16.0.0 replaces the deprecated Elasticsearch High Level REST Client with the Elasticsearch Java API Client. Elastic supports this client against Elasticsearch 8.19.x and later 8.x servers, and against 9.x servers through REST API compatibility, which the connector enables automatically. Before you upgrade from connector version 15.x to 16.x, note the following changes:

  • Connector version 16.x does not support Elasticsearch 7.x or earlier. Against these servers, connector configuration validation fails with an error that begins Elasticsearch version X is not supported by connector version Y. Use the connector version 15.x for these servers.

  • No configuration properties were added, removed, or renamed, and no defaults changed. The max.in.flight.requests property behaves differently. In connector version 15.x, the concurrent bulk requests were one less than the configured value, so max.in.flight.requests=2 sent one request at a time. Connector version 16.x sends as many concurrent requests as the configured number. Inserts with key.ignore=false are protected by offset-based versioning at any value. Upserts carry no version, so two updates to the same document ID can now reach Elasticsearch out of order. See the following warning and max.in.flight.requests.

  • The connector now requires the X-Elastic-Product: Elasticsearch response header from the server. See Limitations.

  • The plugin name stays unchanged. Install version 16.x on each Connect worker and restart the connector. The connector reuses the committed offsets. Existing indices, data streams, aliases, and mappings remain unaffected.

Warning

Upgrading can change write ordering for write.method=upsert. If you relied on ordering with max.in.flight.requests=2 in version 15.x, set max.in.flight.requests=1 before you upgrade. If you used 1, continue using it. For the same concurrency as version 15.x at other values, decrease the value by one. The connector does not warn about this change at startup.

Install the Elasticsearch Service Sink Connector

You can install this connector by using the confluent connect plugin install command, or by manually downloading the ZIP file.

Prerequisites

  • You must install the connector on every machine where Connect will run.

  • Kafka Broker: Confluent Platform 7.9.0 or later, or Kafka 3.9.0 or later.

  • Connect: Confluent Platform 7.9.0 or later, or Kafka 3.9.0 or later.

  • Java 1.8.

  • Elasticsearch 8.x or 9.x.

  • Elastic’s client compatibility guarantee covers only Elasticsearch 8.19.x and later, so the recommended versions are Elasticsearch 8.19.x or 9.x. For Elasticsearch 7.x or earlier, use the 15.x version of this connector.

  • Elasticsearch assigned privileges: create_index, read, write, view_index_metadata on the indices the connector writes to, and the cluster monitor privilege so that the connector can read the server version.

    The following example grants the privileges on indices that match the test-elasticsearch-sink* string. Replace the pattern with the one that matches your indices.

    curl -XPOST "localhost:9200/_security/role/es_sink_connector_role?pretty" -H 'Content-Type: application/json' -d'
    {
       "cluster": [ "monitor" ],
       "indices": [
          {
             "names": [ "test-elasticsearch-sink*" ],
             "privileges": ["create_index", "read", "write", "view_index_metadata"]
          }
       ]
    }'
    
    curl -XPOST "localhost:9200/_security/user/es_sink_connector_user?pretty" -H 'Content-Type: application/json' -d'
    {
    "password" : "seCret-secUre-PaSsW0rD",
    "roles" : [ "es_sink_connector_role" ]
    }'
    
  • An installation of the latest (latest) connector version.

    To install the latest connector version, navigate to your Confluent Platform installation directory and run the following command:

    confluent connect plugin install confluentinc/kafka-connect-elasticsearch:latest
    

    You can install a specific version by replacing latest with a version number as shown in the following example:

    confluent connect plugin install confluentinc/kafka-connect-elasticsearch:16.0.0
    

Install the connector manually

Download and extract the ZIP file for your connector and then follow the manual connector installation instructions.

Quick Start

This quick start uses the Elasticsearch connector to export data produced by the Avro console producer to Elasticsearch.

Prerequisites

See also

For a more detailed Docker-based example of the Confluent Elasticsearch Connector, refer to Confluent Platform Demo (cp-demo). You can deploy a Kafka streaming ETL, including Elasticsearch, using ksqlDB for stream processing.

This quick start assumes that you are using the Confluent CLI commands, but standalone installations are also supported. By default ZooKeeper, Apache Kafka®, Schema Registry, Kafka Connect REST API, and Kafka Connect are started with the confluent local start command. Note that as of Confluent Platform 7.5, ZooKeeper is deprecated for new deployments. Confluent recommends KRaft mode for new deployments.

Add a record to the consumer

  1. Start the Avro console producer to import a few records to Kafka:

    ${CONFLUENT_HOME}/bin/kafka-avro-console-producer \
    --broker-list localhost:9092 --topic test-elasticsearch-sink \
    --property value.schema='{"type":"record","name":"myrecord","fields":[{"name":"f1","type":"string"}]}'
    
  2. Enter the following in the console producer:

    {"f1": "value1"}
    {"f1": "value2"}
    {"f1": "value3"}
    

    The three records entered are published to the Kafka topic test-elasticsearch in Avro format.

Load the connector

Complete the following steps to load the predefined Elasticsearch connector bundled with Confluent Platform.

Note

Default connector properties are already set for this quick start. To view the connector properties, refer to etc/kafka-connect-elasticsearch/quickstart-elasticsearch.properties.

  1. List the available predefined connectors using the following command:

    confluent local list
    

    Example output:

    Bundled Predefined Connectors (edit configuration under etc/):
      elasticsearch-sink
      file-source
      file-sink
      jdbc-source
      jdbc-sink
      hdfs-sink
      s3-sink
    
  2. Load the elasticsearch-sink connector:

    confluent local load elasticsearch-sink
    

    Example output:

    {
      "name": "elasticsearch-sink",
      "config": {
        "connector.class": "io.confluent.connect.elasticsearch.ElasticsearchSinkConnector",
        "tasks.max": "1",
        "topics": "test-elasticsearch-sink",
        "key.ignore": "true",
        "connection.url": "http://localhost:9200",
        "name": "elasticsearch-sink"
      },
      "tasks": [],
      "type": null
    }
    

    Tip

    For non-CLI users, you can load the Elasticsearch connector by running Kafka Connect in standalone mode with this command:

    ./bin/connect-standalone etc/schema-registry/connect-avro-standalone.properties \
    etc/kafka-connect-elasticsearch/quickstart-elasticsearch.properties
    
  3. After the connector finishes ingesting data to Elasticsearch, enter the following command to check that data is available in Elasticsearch:

    curl -XGET 'http://localhost:9200/test-elasticsearch-sink/_search?pretty'
    

    Example output:

    {
      "took" : 39,
      "timed_out" : false,
      "_shards" : {
        "total" : 1,
        "successful" : 1,
        "skipped" : 0,
        "failed" : 0
      },
      "hits" : {
        "total" : {
          "value" : 3,
          "relation" : "eq"
        },
        "max_score" : 1.0,
        "hits" : [
          {
            "_index" : "test-elasticsearch-sink",
            "_id" : "test-elasticsearch-sink+0+0",
            "_score" : 1.0,
            "_source" : {
              "f1" : "value1"
            }
          },
          {
            "_index" : "test-elasticsearch-sink",
            "_id" : "test-elasticsearch-sink+0+2",
            "_score" : 1.0,
            "_source" : {
              "f1" : "value3"
            }
          },
          {
            "_index" : "test-elasticsearch-sink",
            "_id" : "test-elasticsearch-sink+0+1",
            "_score" : 1.0,
            "_source" : {
              "f1" : "value2"
            }
          }
        ]
      }
    }
    

Delivery semantics

The connector delivers records at least once and relies on Elasticsearch’s idempotent write semantics to avoid duplicate documents. Inserts and deletes with key.ignore=false additionally use external versioning so that each key converges on its latest record (see Delivery guarantees). To boost throughput, it also batches and pipelines writes, accumulating messages into batches and processing batches concurrently. The number of concurrent bulk requests is controlled by max.in.flight.requests.

Mapping management

Mapping determines how Elasticsearch tokenizes, analyzes, and indexes your data, and you can’t change some mapping properties after they’re defined. You can add new fields to an index, but you can’t add new analyzers or change existing fields. Changing an existing field after data is indexed makes the already-indexed data incorrect and breaks your searches. Define mappings before you write data to Elasticsearch.

Index templates can be helpful when manually defining mappings, and allow you to define templates that are automatically applied when new indices are created. The templates include both settings and mappings, along with a simple pattern template that controls whether the template should be applied to the new index.

Schema evolution

Each Kafka topic writes to its own Elasticsearch index, so schemas for different topics evolve independently. Elasticsearch enforces only one constraint on mappings—all fields with the same name in the same index must have the same mapping—which is what makes this independent evolution possible.

The connector infers a mapping only once per index, data stream, or alias. When schema.ignore is false, it builds the mapping from the schema of the first record written to that resource and creates the mapping if the resource does not have one. The connector does not update an existing mapping when the record schema changes later. Any schema evolution after that point depends on Elasticsearch dynamic mapping, described next. When schema.ignore is true, the connector never creates a mapping and Elasticsearch infers every field dynamically.

Elasticsearch supports dynamic mapping: when it encounters previously unknown field in a document, it uses dynamic mapping to determine the datatype for the field and automatically adds the new field to the index mapping.

When dynamic mapping is enabled, Elasticsearch absorbs schema changes in the records the connector writes. This is because mappings in Elasticsearch are more flexible than the schema evolution allowed in Connect when different converters are used. For example, when the Avro converter is used, backward, forward, and fully compatible schema evolutions are allowed.

When dynamic mapping is enabled, Elasticsearch accepts the following schema changes in records written by the connector:

  • Adding Fields: Adding one or more fields to Kafka messages. Elasticsearch adds the new fields to the mapping when dynamic mapping is enabled.

  • Removing Fields: Removing one or more fields from Kafka messages. Missing fields are treated as the null value defined for those fields in the mapping.

  • Changing types that can be merged: Changing a field from integer type to string type. Elasticsearch can convert integers to strings.

The following change is not allowed:

  • Changing types that can not be merged: Changing a field from a string type to an integer type.

Because mappings are more flexible, schema compatibility should be enforced when writing data to Kafka.

Automatic retries

The connector automatically retries a failed bulk request five times by default before marking the task as failed, using an exponential backoff technique to give the Elasticsearch service time to recover when it’s temporarily overloaded. This technique adds randomness, called jitter, to the calculated backoff times to prevent a thundering herd, wherein large numbers of requests from many tasks are submitted concurrently and overwhelm the service.

Retries apply to failures of a request as a whole: connection errors, timeouts, and error HTTP responses to the entire request, including an HTTP 429 returned for the whole request. The same retry settings apply to the requests that create indices, data streams, and mappings. The connector never retries individual documents. When Elasticsearch rejects a document inside an otherwise successful bulk response, rejections caused by mapping or parsing errors follow the behavior.on.malformed.documents property. Any other per-document rejection, including a per-document 429, is reported to the dead letter queue, if one is configured, and then fails the task.

Randomness spreads out the retries from many tasks. This should reduce the overall time required to complete all outstanding requests when compared to simple exponential backoff. The goal is to spread out the requests to Elasticsearch as much as possible.

The number of retries is dictated by the max.retries connector configuration property. The max.retries property defaults to five attempts. The maximum backoff time (the amount of time to wait before retrying) is a function of the retry attempt number and the initial backoff time specified in the retry.backoff.ms connector configuration property. The retry.backoff.ms property defaults to 100 milliseconds.

The jitter strategy used is “Full Jitter” where the actual backoff time is a uniform random value selected between the minimum backoff (0.0) and maximum backoff at the current attempt. Since the actual backoff value is selected randomly, it is not guaranteed to increase with each consecutive retry attempt.

For example, the following table shows the possible wait times for four subsequent retries with retry.backoff.ms set to 500 milliseconds (0.5 second):

Range of backoff times

Retry

Minimum Backoff (sec)

Maximum Backoff (sec)

Actual Backoff with Jitter (sec)

Total Potential Delay from First Attempt (sec)

1

0.0

1.0

0.4

1.0

2

0.0

2.0

1.4

3.0

3

0.0

4.0

3.8

7.0

4

0.0

8.0

3.0

15.0

Note how the maximum wait time is the normal exponential backoff, which is calculated as ${retry.backoff.ms} * 2 ^ retry, so the first retry waits at most twice retry.backoff.ms. Also note how the actual backoff decreased between retry attempt #3 and #4 despite the maximum backoff increasing exponentially. A single backoff is capped at 24 hours.

As shown in the following table, increasing the maximum number of retries adds more backoff:

Range of backoff times for additional retries

Retry

Minimum Backoff (sec)

Maximum Backoff (sec)

Total Potential Delay from First Attempt (sec)

5

0.0

16.0

31.0

6

0.0

32.0

63.0

7

0.0

64.0

127.0

8

0.0

128.0

255.0

9

0.0

256.0

511.0

10

0.0

512.0

1,023.0

11

0.0

1,024.0

2,047.0

12

0.0

2,048.0

4,095.0

13

0.0

4,096.0

8,191.0

By increasing max.retries to 10, the connector can take up to 1,023 seconds, or a little over 17 minutes, to successfully send a batch of records when the Elasticsearch service is overloaded. Increasing the value to 13 quickly increases the maximum potential time to submit a batch of records to well over two hours.

You can adjust both the max.retries and retry.backoff.ms connector configuration properties to optimize retry timing.

Reindexing

Reindexing lets you change how a set of documents is indexed—for example, updating the analyzer, tokenizer, or indexed fields—even though these properties can’t be changed on a mapping that’s already defined. You can use Index aliases to achieve reindexing with zero downtime.

To reindex the data, complete the following steps in Elasticsearch:

  1. Create an alias for the index with the original mapping.

  2. Point the applications using the index to the alias.

  3. Create a new index with the updated mapping.

  4. Move data from the original index to the new index.

  5. Atomically move the alias to the new index.

  6. Delete the original index.

Write requests continue to come during the reindex period (if reindexing is done with no downtime). Aliases do not allow writing to both the original and the new index at the same time. To solve this, you can use two Elasticsearch connector jobs to achieve double writes, one to the original index and a second one to the new index. The following steps explain how to do this:

  1. Keep the original connector job that ingests data to the original indices running.

  2. Create a new connector job that writes to new indices. As long as the data is in Kafka, some of the old data and all new data are written to the new indices.

  3. After the reindexing process is complete and the data in the original indices are moved to the new indices, stop the original connector job.

Security

The Elasticsearch connector can read data from secure Kafka by following the instructions in the Kafka Connect security documentation. The connector can write data to a secure Elasticsearch cluster that supports either of the following authentication methods:

  • Basic authentication: By setting the connection.username and connection.password configuration properties.

  • Kerberos authentication: By setting the Kerberos configuration properties.

For more information, see Elasticsearch Connector with Security.

Suggested Resources