Esql

Esql

Certified

Run ES|QL query

yaml
type: io.kestra.plugin.elasticsearch.Esql
yaml
id: esql_error_report
namespace: company.team

inputs:
  - id: service
    type: STRING
    defaults: "api"
  - id: since
    type: STRING
    defaults: "2024-01-01T00:00:00Z"

tasks:
  - id: query
    type: io.kestra.plugin.elasticsearch.Esql
    fetchType: STORE
    async: true
    query: |
      FROM logs
      | WHERE service == ? AND @timestamp >= ?
      | STATS error_count = COUNT(*) BY level, host
      | SORT error_count DESC
    params:
      - "{{ inputs.service }}"
      - "{{ inputs.since }}"

pluginDefaults:
  - type: io.kestra.plugin.elasticsearch
    values:
      connection:
        headers:
          - "Authorization: ApiKey yourEncodedApiKey"
        hosts:
          - https://yourCluster.us-central1.gcp.cloud.es.io:443

yaml
id: bulk_load_and_query
namespace: company.team

tasks:
  - id: extract
    type: io.kestra.plugin.core.http.Download
    uri: https://huggingface.co/datasets/kestra/datasets/resolve/main/jsonl/books.jsonl

  - id: load
    type: io.kestra.plugin.elasticsearch.Bulk
    from: "{{ outputs.extract.uri }}"

  - id: sleep
    type: io.kestra.plugin.core.flow.Sleep
    duration: PT5S
    description: Pause needed after load before we can query

  - id: query
    type: io.kestra.plugin.elasticsearch.Esql
    fetchType: STORE
    query: |
      FROM books
        | KEEP author, name, page_count, release_date
        | SORT page_count DESC
        | LIMIT 5

pluginDefaults:
  - type: io.kestra.plugin.elasticsearch
    values:
      connection:
        headers:
          - "Authorization: ApiKey yourEncodedApiKey"
        hosts:
          - https://yourCluster.us-central1.gcp.cloud.es.io:443
Properties
Definitions
hosts*Requiredarray
SubTypestring
Min items1
basicAuth
passwordstring
usernamestring
headersarray
SubTypestring
pathPrefixstring
strictDeprecationModebooleanstring
targetServerVersionintegerstring
Default8
trustAllSslbooleanstring
Defaultfalse
Defaultfalse
DefaultFETCH
Possible Values
STOREFETCHFETCH_ONENONE
SubTypestring
SubTypeobject
Formaturi
Unitrecords