Elasticsearch Connector with Security

Complete the following instructions to connect the Kafka Connect Elasticsearch Sink connector to an Elasticsearch cluster over TLS with username and password authentication. The instructions follow the Elasticsearch guide Set up basic security plus HTTPS and were verified against Elasticsearch 8.19. They also apply to Elasticsearch 9.x.

In Elasticsearch 8.x and later, security is enabled by default: the HTTP layer requires TLS, and every request must be authenticated. The connector verifies the server certificate against a truststore that you configure with the elastic.https.ssl.* properties and authenticates with connection.username and connection.password. Optionally, you can also present a client certificate (mutual TLS).

Prerequisites:

Note

This walkthrough configures TLS before Elasticsearch starts for the first time, so that you control the certificates. If you have already started Elasticsearch once, it has generated its own certificates. In that case, skip to Use the certificates that Elasticsearch generated.

Step 1: Download and extract the Elasticsearch archive

Download the archive for your platform and extract it. Do not start Elasticsearch yet. The examples below refer to the extracted directory as $ES_HOME.

tar xzvf elasticsearch-<version>-<platform>.tar.gz
cd elasticsearch-<version>
export ES_HOME=$(pwd)

Step 2: Generate certificates

Elasticsearch ships a certificate utility, elasticsearch-certutil, that creates a certificate authority (CA) and certificates signed by it. Use it to generate a CA and a server certificate whose subject alternative names match the host name in the connector’s connection.url.

  1. Generate the CA:

    mkdir -p $ES_HOME/config/certs
    $ES_HOME/bin/elasticsearch-certutil ca --pem --out $ES_HOME/config/certs/ca.zip
    unzip $ES_HOME/config/certs/ca.zip -d $ES_HOME/config/certs
    

    This creates config/certs/ca/ca.crt and config/certs/ca/ca.key.

  2. Describe the server certificate. The dns and ip entries must include every name the connector uses to reach Elasticsearch. The connector verifies the host name by default.

    cat <<EOF > $ES_HOME/config/certs/instances.yml
    instances:
      - name: elasticsearch
        dns:
          - localhost
        ip:
          - 127.0.0.1
    EOF
    
  3. Generate the server certificate:

    $ES_HOME/bin/elasticsearch-certutil cert --pem \
      --in $ES_HOME/config/certs/instances.yml \
      --ca-cert $ES_HOME/config/certs/ca/ca.crt \
      --ca-key $ES_HOME/config/certs/ca/ca.key \
      --out $ES_HOME/config/certs/certs.zip
    unzip $ES_HOME/config/certs/certs.zip -d $ES_HOME/config/certs
    

    This creates config/certs/elasticsearch/elasticsearch.crt and config/certs/elasticsearch/elasticsearch.key.

  4. (Optional, for mutual TLS) Generate a client certificate for the connector:

    $ES_HOME/bin/elasticsearch-certutil cert --pem --name connector \
      --ca-cert $ES_HOME/config/certs/ca/ca.crt \
      --ca-key $ES_HOME/config/certs/ca/ca.key \
      --out $ES_HOME/config/certs/connector.zip
    unzip $ES_HOME/config/certs/connector.zip -d $ES_HOME/config/certs
    

    This creates config/certs/connector/connector.crt and config/certs/connector/connector.key.

Step 3: Configure and start Elasticsearch

  1. Append the TLS settings to config/elasticsearch.yml. Set xpack.security.http.ssl.client_authentication to required only if you want the connector to present a client certificate.

    cat <<EOF >> $ES_HOME/config/elasticsearch.yml
    xpack.security.enabled: true
    xpack.security.http.ssl.enabled: true
    xpack.security.http.ssl.key: certs/elasticsearch/elasticsearch.key
    xpack.security.http.ssl.certificate: certs/elasticsearch/elasticsearch.crt
    xpack.security.http.ssl.certificate_authorities: [ "certs/ca/ca.crt" ]
    xpack.security.http.ssl.client_authentication: optional
    xpack.security.transport.ssl.enabled: true
    xpack.security.transport.ssl.key: certs/elasticsearch/elasticsearch.key
    xpack.security.transport.ssl.certificate: certs/elasticsearch/elasticsearch.crt
    xpack.security.transport.ssl.certificate_authorities: [ "certs/ca/ca.crt" ]
    EOF
    

    Transport TLS is not used by the connector, but Elasticsearch requires it whenever security is enabled on a production node.

    If you set client_authentication to required, every client must present a certificate signed by the CA, including the curl commands in this walkthrough. Add --cert $ES_HOME/config/certs/connector/connector.crt --key $ES_HOME/config/certs/connector/connector.key to each curl command below.

  2. Set the password for the built-in elastic superuser. Elasticsearch reads it from the bootstrap.password entry in its keystore on first start:

    $ES_HOME/bin/elasticsearch-keystore create
    $ES_HOME/bin/elasticsearch-keystore add bootstrap.password
    
  3. Start Elasticsearch:

    $ES_HOME/bin/elasticsearch
    
  4. In a new terminal, test the connection. curl verifies the server certificate against the CA and authenticates as elastic:

    curl --cacert $ES_HOME/config/certs/ca/ca.crt -u elastic \
      https://localhost:9200
    
  5. Create a role and a user for the connector. The connector needs the create_index, read, write, and view_index_metadata privileges on the indices it writes to, and the cluster monitor privilege to read the server version at startup. Replace the names pattern with one that matches your indices.

    curl --cacert $ES_HOME/config/certs/ca/ca.crt -u elastic \
      -XPOST "https://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 --cacert $ES_HOME/config/certs/ca/ca.crt -u elastic \
      -XPOST "https://localhost:9200/_security/user/es_sink_connector_user?pretty" \
      -H 'Content-Type: application/json' -d'
    {
      "password" : "seCret-secUre-PaSsW0rD",
      "roles" : [ "es_sink_connector_role" ]
    }'
    

    Without the monitor privilege the connector still works, but it logs Failed to get ES server version at startup and cannot validate the Elasticsearch version.

Step 4: Create the connector truststore

The connector trusts the server certificate through a Java truststore that contains the CA certificate. Choose a password when prompted.

keytool -importcert -noprompt -alias elasticsearch-ca \
  -file $ES_HOME/config/certs/ca/ca.crt \
  -keystore $ES_HOME/config/certs/truststore.jks

(Optional, for mutual TLS) Package the connector’s client certificate and key as a PKCS12 keystore:

openssl pkcs12 -export -name connector \
  -in $ES_HOME/config/certs/connector/connector.crt \
  -inkey $ES_HOME/config/certs/connector/connector.key \
  -out $ES_HOME/config/certs/keystore.p12

Step 5: Configure the connector

Prerequisites
  1. Open a new terminal and change your current directory to <path-to-confluent>.

  2. Save the following configuration as etc/kafka-connect-elasticsearch/elastic-secure.properties. Replace the paths and passwords with your own. The elastic.security.protocol=SSL property enables the elastic.https.ssl.* properties.

    cat <<EOF > etc/kafka-connect-elasticsearch/elastic-secure.properties
    name=elasticsearch-sink-secure
    connector.class=io.confluent.connect.elasticsearch.ElasticsearchSinkConnector
    tasks.max=1
    topics=test-elasticsearch-sink
    key.ignore=true
    connection.url=https://localhost:9200
    connection.username=es_sink_connector_user
    connection.password=seCret-secUre-PaSsW0rD
    
    elastic.security.protocol=SSL
    elastic.https.ssl.truststore.location=$ES_HOME/config/certs/truststore.jks
    elastic.https.ssl.truststore.password=<truststore-password>
    elastic.https.ssl.truststore.type=JKS
    EOF
    

    For mutual TLS, also add the keystore properties and set xpack.security.http.ssl.client_authentication to required in Step 3:

    elastic.https.ssl.keystore.location=$ES_HOME/config/certs/keystore.p12
    elastic.https.ssl.keystore.password=<keystore-password>
    elastic.https.ssl.keystore.type=PKCS12
    

    Keep connection.username and connection.password. A client certificate satisfies the TLS layer but does not authenticate the user unless you configure a PKI realm in Elasticsearch. Without credentials, Elasticsearch returns 401. Without the keystore, connector configuration validation fails with Received fatal alert: certificate_required.

    Leave elastic.https.ssl.protocol and elastic.https.ssl.enabled.protocols at their defaults. For the full list of TLS properties, see the Security configuration properties.

  3. Start Connect and load the connector:

    bin/confluent local services connect start
    
    bin/confluent local services connect connector load elasticsearch-sink-secure \
      --config etc/kafka-connect-elasticsearch/elastic-secure.properties
    

Step 6: Test the system

  1. Produce a few records:

    bin/kafka-avro-console-producer --bootstrap-server localhost:9092 \
      --topic test-elasticsearch-sink \
      --property value.schema='{"type":"record","name":"myrecord","fields":[{"name":"f1","type":"string"}]}'
    {"f1": "secret1"}
    {"f1": "secret2"}
    
  2. Query Elasticsearch as the connector user:

    curl --cacert $ES_HOME/config/certs/ca/ca.crt \
      -u es_sink_connector_user:seCret-secUre-PaSsW0rD \
      'https://localhost:9200/test-elasticsearch-sink/_search?pretty'
    

Use the certificates that Elasticsearch generated

If you started Elasticsearch 8.x or 9.x without configuring TLS, it enabled security automatically on first start, generated a CA and server certificate, and printed the elastic password. See Security auto-configuration. To connect to such a node:

  1. Build the truststore from the generated CA certificate, $ES_HOME/config/certs/http_ca.crt:

    keytool -importcert -noprompt -alias elasticsearch-ca \
      -file $ES_HOME/config/certs/http_ca.crt \
      -keystore $ES_HOME/config/certs/truststore.jks
    
  2. If you no longer have the elastic password, reset it:

    $ES_HOME/bin/elasticsearch-reset-password -u elastic
    
  3. Create the connector role and user as in Step 3, and configure the connector as in Step 5. The generated server certificate includes the node’s host name, localhost, and its IP addresses, so host name verification works for local connections.

Troubleshooting

  • Host name verification fails: the host in connection.url is not in the server certificate’s subject alternative names. Add it to instances.yml and regenerate the certificate. Do not turn off verification with elastic.https.ssl.endpoint.identification.algorithm= unless you cannot add the host name to the certificate.

  • Connector validation fails with certificate_required: Elasticsearch has client_authentication: required and the connector has no elastic.https.ssl.keystore.location. Add the keystore properties or set client_authentication to optional.

  • Elasticsearch returns 401 although you configured a client certificate: a client certificate alone does not authenticate a user. Set connection.username and connection.password, and make sure the user exists. If client_authentication is required, the curl commands that create the role and user must also present the client certificate.

  • Elasticsearch returns 403 on GET /: the connector user lacks the cluster monitor privilege. The connector logs a warning and continues.

  • Connector configuration validation fails with Missing [X-Elastic-Product] header: a proxy or load balancer between the connector and Elasticsearch removes the header. See Limitations.