Spec-Zone.ru › Elasticsearch 8
›Руководство по Elasticsearch [8.17] ›REST API ›API документов

API обновления по запросу

Новая справочная информация по API

Для получения самой актуальной информации об API обратитесь к API документов.

Обновляет документы, соответствующие заданному запросу. Если запрос не указан, выполняет обновление всех документов в потоке данных или индексе без изменения исходных данных, что полезно для подбора изменений в схеме.

resp = client.update_by_query(
    index="my-index-000001",
    conflicts="proceed",
)
print(resp)
response = client.update_by_query(
  index: 'my-index-000001',
  conflicts: 'proceed'
)
puts response
const response = await client.updateByQuery({
  index: "my-index-000001",
  conflicts: "proceed",
});
console.log(response);
POST my-index-000001/_update_by_query?conflicts=proceed

Запрос

POST /<target>/_update_by_query

Предварительные требования

  • Если включены функции безопасности Elasticsearch, вам необходимо иметь следующие права доступа к индексам для целевого потока данных, индекса или псевдонима:

    • read
    • index или write

Описание

Вы можете указать критерии запроса в URI запроса или теле запроса, используя тот же синтаксис, что и в API поиска.

При отправке запроса обновления по запросу Elasticsearch получает моментальный снимок потока данных или индекса в момент начала обработки запроса и обновляет соответствующие документы с использованием internal версиирования. Если версии совпадают, документ обновляется, а номер версии увеличивается. Если документ изменяется между моментом создания моментального снимка и обработкой операции обновления, это приводит к конфликту версий, и операция завершается неудачно. Вы можете выбрать подсчет конфликтов версий вместо остановки и возврата, установив conflicts в proceed. Обратите внимание, что если вы выберете подсчет конфликтов версий, операция может попытаться обновить больше документов из источника, чем max_docs, пока не успешно обновит max_docs документов или не пройдёт по всем документам в исходном запросе.

Документы с версией, равной 0, не могут быть обновлены с помощью обновления по запросу, поскольку internal версия не поддерживает 0 в качестве допустимого номера версии.

При обработке запроса обновления по запросу Elasticsearch выполняет несколько поисковых запросов последовательно, чтобы найти все соответствующие документы. Для каждой группы соответствующих документов выполняется пакетное обновление. Любые сбои в запросе или обновлении приводят к ошибке запроса обновления по запросу, и ошибки отображаются в ответе. Все успешно завершенные запросы обновлений остаются применёнными, они не отменяются.

Обновление фрагментов

Указание параметра refresh обновляет все фрагменты после завершения запроса. Это отличается от параметра refresh API обновления, который обновляет только фрагмент, получивший запрос. В отличие от API обновления, он не поддерживает wait_for.

Асинхронное выполнение запросов обновления по запросу

Если запрос содержит wait_for_completion=false, Elasticsearch выполняет некоторые предварительные проверки, запускает запрос и возвращает task, который можно использовать для отмены или получения состояния задачи. Elasticsearch создаёт запись этой задачи в виде документа в .tasks/task/${taskId}.

Ожидание активных фрагментов

wait_for_active_shards управляет тем, сколько копий фрагмента должно быть активными перед продолжением запроса. Подробности см. в разделе Активные фрагменты. timeout управляет тем, как долго каждый запрос записи ожидает, пока недоступные фрагменты не станут доступными. Оба работают точно так же, как и в API пакетной обработки. Обновление по запросу использует прокручиваемые запросы, поэтому вы также можете указать параметр scroll для управления временем сохранения контекста поиска, например ?scroll=10m. По умолчанию — 5 минут.

Ограничение запросов обновления

Для управления скоростью, с которой обновление по запросу отправляет пакеты операций обновления, вы можете установить requests_per_second на любое положительное десятичное число. Это добавляет время ожидания к каждому пакету для ограничения скорости. Установите requests_per_second в -1, чтобы отключить ограничение.

Ограничение использует время ожидания между пакетами, чтобы внутренние запросы прокрутки могли получить таймаут, учитывающий добавочное время запроса. Время добавления — это разница между размером пакета, делённого на requests_per_second, и временем записи. По умолчанию размер пакета составляет 1000, поэтому, если requests_per_second установлено в 500:

target_time = 1000 / 500 per second = 2 seconds
wait_time = target_time - write_time = 2 seconds - .5 seconds = 1.5 seconds

Поскольку пакет отправляется как один _bulk запрос, большие размеры пакетов заставляют Elasticsearch создавать множество запросов и ожидать перед началом следующей группы. Это «импульсный», а не «плавный» режим.

Разбиение

Обновление по запросу поддерживает прокрутку с нарезкой для параллелизации процесса обновления. Это может повысить эффективность и предоставить удобный способ разбить запрос на меньшие части.

Установка slices в auto выбирает разумное значение для большинства потоков данных и индексов. Если вы выполняете разбиение вручную или настраиваете автоматическое разбиение, имейте в виду, что:

  • Производительность запроса наилучшим образом достигается, когда количество slices равно количеству фрагментов в индексе или базовом индексе. Если это число большое (например, 500), выберите меньшее значение, так как слишком большое количество slices ухудшает производительность. Увеличение slices значения сверх количества фрагментов обычно не улучшает эффективность и добавляет накладные расходы.
  • Производительность обновления линейно масштабируется с помощью доступных ресурсов с количеством срезов.

Производительность запроса или обновления доминирует над временем выполнения в зависимости от переиндексируемых документов и ресурсов кластера.

Параметры пути

<target>
(Необязательный, строка) Список потоков данных, индексов и псевдонимов для поиска, разделённых запятыми. Поддерживаются подстановочные знаки (*). Для поиска во всех потоках данных или индексах опустите этот параметр или используйте * или _all.

Параметры запроса

allow_no_indices

(Необязательно, логическое значение) Если false, запрос вернёт ошибку, если какой-либо шаблонный запрос, псевдоним индекса или _all значение указывают только на отсутствующие или закрытые индексы. Это поведение применяется даже если запрос нацелен на другие открытые индексы. Например, запрос, нацеленный на foo*,bar*, возвращает ошибку, если индекс начинается с foo, но ни один индекс не начинается с bar.

По умолчанию true.

analyzer

(Необязательно, строка) Анализатор, используемый для строки запроса.

Этот параметр может быть использован только при указании параметра строки запроса q.

analyze_wildcard

(Необязательно, логическое значение) Если true, выполняются анализ запросов с подстановкой и префиксом. По умолчанию false.

Этот параметр может быть использован только при указании параметра строки запроса q.

conflicts
(Необязательно, строка) Действие при конфликте версий при обновлении по запросу: abort или proceed. По умолчанию abort.
default_operator

(Необязательно, строка) Операторы по умолчанию для запросов в формате строки запроса: AND или OR. По умолчанию OR.

Этот параметр может быть использован только при указании параметра строки запроса q.

df

(Необязательно, строка) Поле, используемое по умолчанию, если в строке запроса нет префикса поля.

Этот параметр может быть использован только при указании параметра строки запроса q.

expand_wildcards

(Необязательно, строка) Тип индекса, с которым могут соответствовать шаблоны подстановки. Если запрос может нацеливаться на потоки данных, этот аргумент определяет, соответствуют ли шаблоны подстановки скрытым потокам данных. Поддерживает значения, разделенные запятыми, такие как open,hidden. Допустимые значения:

all
Сопоставление любого потока данных или индекса, включая скрытые.
open
Сопоставление открытых, нескрытых индексов. Также соответствует любым нескрытым потокам данных.
closed
Сопоставление закрытых, нескрытых индексов. Также соответствует любым нескрытым потокам данных. Потоки данных не могут быть закрыты.
hidden
Сопоставление скрытых потоков данных и скрытых индексов. Должно быть использовано совместно с open, closed или обоими.
none
Шаблоны подстановки не принимаются.

По умолчанию open.

ignore_unavailable
(Необязательно, логическое значение) Если false, запрос вернёт ошибку, если он нацелен на отсутствующий или закрытый индекс. По умолчанию false.
lenient

(Необязательно, логическое значение) Если true, ошибки запроса, основанные на формате (например, предоставление текста числовому полю) в строке запроса будут проигнорированы. По умолчанию false.

Этот параметр может быть использован только при указании параметра строки запроса q.

max_docs
(Необязательно, целое число) Максимальное количество документов для обработки. По умолчанию все документы. Если значение меньше или равно scroll_size, для получения результатов операции не будет использоваться скроллинг.
pipeline
(Необязательно, строка) Идентификатор конвейера для предварительной обработки входящих документов. Если в индексе указан конвейер по умолчанию, установка значения на _none отключит конвейер по умолчанию для этого запроса. Если настроен конечный конвейер, он всегда будет выполнен, независимо от значения этого параметра.
preference
(Необязательно, строка) Указывает узел или фрагмент, на котором должна быть выполнена операция. По умолчанию случайный.
q
(Необязательно, строка) Запрос в синтаксисе строки запроса Lucene.
request_cache
(Необязательно, логическое значение) Если true, кэширование запросов используется для этого запроса. По умолчанию используется настройка на уровне индекса.
refresh
(Необязательно, логическое значение) Если true, Elasticsearch обновляет затронутые фрагменты, чтобы сделать операцию видимой для поиска. По умолчанию false.
requests_per_second
(Необязательно, целое число) Скорость этого запроса в подзапросах в секунду. По умолчанию -1 (без ограничения скорости).
routing
(Необязательно, строка) Пользовательское значение, используемое для маршрутизации операций на определённый фрагмент.
scroll
(Необязательно, значение времени) Период сохранения контекста поиска для скроллинга. См. Скроллинг результатов поиска.
scroll_size
(Необязательно, целое число) Размер запроса скроллинга, который питает операцию. По умолчанию 1000.
search_type

(Необязательно, строка) Тип операции поиска. Доступные варианты:

  • query_then_fetch
  • dfs_query_then_fetch
search_timeout
(Необязательно, единицы времени) Явное время ожидания для каждого запроса поиска. По умолчанию без ограничения времени.
slices
(Необязательно, целое число) Количество частей, на которые следует разделить эту задачу. По умолчанию 1, что означает, что задача не разбивается на подзадачи.
sort
(Необязательно, строка) Список пар <поле>:<направление>, разделённых запятыми.
stats
(Необязательно, строка) Конкретный tag запроса для ведения журнала и статистических целей.
terminate_after

(Необязательно, целое число) Максимальное количество документов для сбора для каждого фрагмента. Если запрос достигает этого предела, Elasticsearch прерывает запрос. Elasticsearch собирает документы до сортировки.

Используйте с осторожностью. Elasticsearch применяет этот параметр к каждому фрагменту, обрабатывающему запрос. По возможности позвольте Elasticsearch выполнять раннее прерывание автоматически. Избегайте указания этого параметра для запросов, которые нацелены на потоки данных с поддержкой индексов на нескольких уровнях данных.

timeout

(Необязательно, единицы времени) Период, в течение которого каждый запрос обновления ожидает следующих операций:

  • Динамические обновления отображения
  • Ожидание активных фрагментов

По умолчанию 1m (одна минута). Это гарантирует, что Elasticsearch подождёт хотя бы до времени ожидания, прежде чем завершить работу с ошибкой. Фактическое время ожидания может быть больше, особенно при нескольких ожиданиях.

version
(Необязательно, логическое значение) Если true, возвращает версию документа как часть совпадения.
wait_for_active_shards

(Необязательно, строка) Количество копий каждого фрагмента, которые должны быть активны, прежде чем продолжить операцию. Установите на all или любое неотрицательное целое число до общего количества копий каждого фрагмента в индексе (number_of_replicas+1). По умолчанию 1, что означает ожидание только активации каждого первичного фрагмента.

См. Активные фрагменты.

Тело запроса

query
(Необязательно, объект запроса) Указывает документы для обновления с помощью Query DSL.

Тело ответа

took
Количество миллисекунд с начала до конца всей операции.
timed_out
Этот флаг устанавливается в значение true, если какой-либо из запросов, выполненных во время обновления по запросу, истек по времени.
total
Количество документов, которые были обработаны успешно.
updated
Количество документов, которые были обновлены успешно.
deleted
Количество документов, которые были удалены успешно.
batches
Количество ответов скролла, полученных обновлением по запросу.
version_conflicts
Количество конфликтов версий, с которыми столкнулось обновление по запросу.
noops
Количество документов, которые были проигнорированы, потому что скрипт, используемый для обновления по запросу, возвратил значение noop для ctx.op.
retries
Количество попыток повторной обработки, предпринятых обновлением по запросу. bulk — это количество повторно обработанных массовых действий, а search — количество повторно обработанных поисковых действий.
throttled_millis
Количество миллисекунд, в течение которого запрос ожидал, чтобы соответствовать requests_per_second.
requests_per_second
Количество запросов в секунду, фактически выполненных во время обновления по запросу.
throttled_until_millis
Это поле всегда должно быть равно нулю в ответе _update_by_query. Оно имеет смысл только при использовании API задач задач, где оно указывает следующий момент времени (в миллисекундах с начала эпохи), когда запрос, ограниченный по производительности, будет снова выполнен, чтобы соответствовать requests_per_second.
failures
Массив ошибок, если во время процесса произошли невосстановимые ошибки. Если этот массив не пустой, запрос был прерван из-за этих ошибок. Обновление по запросу реализовано с использованием пакетов. Любая ошибка приводит к прерыванию всего процесса, но все ошибки в текущем пакете собираются в массив. Вы можете использовать опцию conflicts, чтобы предотвратить прерывание переиндексации при конфликтах версий.

Примеры

Самое простое использование _update_by_query просто выполняет обновление каждого документа в потоке данных или индексе без изменения источника. Это полезно для получения новой свойства или других изменений онлайн-сопоставления.

Для обновления выбранных документов укажите запрос в теле запроса:

resp = client.update_by_query(
    index="my-index-000001",
    conflicts="proceed",
    query={
        "term": {
            "user.id": "kimchy"
        }
    },
)
print(resp)
response = client.update_by_query(
  index: 'my-index-000001',
  conflicts: 'proceed',
  body: {
    query: {
      term: {
        'user.id' => 'kimchy'
      }
    }
  }
)
puts response
const response = await client.updateByQuery({
  index: "my-index-000001",
  conflicts: "proceed",
  query: {
    term: {
      "user.id": "kimchy",
    },
  },
});
console.log(response);
POST my-index-000001/_update_by_query?conflicts=proceed
{
  "query": { 
    "term": {
      "user.id": "kimchy"
    }
  }
}

Запрос необходимо передать как значение ключу query, так же как и в API поиска. Вы также можете использовать параметр q так же, как и в API поиска.

Обновление документов в нескольких потоках данных или индексах:

resp = client.update_by_query(
    index="my-index-000001,my-index-000002",
)
print(resp)
response = client.update_by_query(
  index: 'my-index-000001,my-index-000002'
)
puts response
const response = await client.updateByQuery({
  index: "my-index-000001,my-index-000002",
});
console.log(response);
POST my-index-000001,my-index-000002/_update_by_query

Ограничить операцию обновления по запросу фрагментами, для которых задано определённое значение маршрутизации:

resp = client.update_by_query(
    index="my-index-000001",
    routing="1",
)
print(resp)
response = client.update_by_query(
  index: 'my-index-000001',
  routing: 1
)
puts response
const response = await client.updateByQuery({
  index: "my-index-000001",
  routing: 1,
});
console.log(response);
POST my-index-000001/_update_by_query?routing=1

По умолчанию обновление по запросу использует пакеты скролла размером 1000. Вы можете изменить размер пакета с помощью параметра scroll_size:

resp = client.update_by_query(
    index="my-index-000001",
    scroll_size="100",
)
print(resp)
response = client.update_by_query(
  index: 'my-index-000001',
  scroll_size: 100
)
puts response
const response = await client.updateByQuery({
  index: "my-index-000001",
  scroll_size: 100,
});
console.log(response);
POST my-index-000001/_update_by_query?scroll_size=100

Обновление документа с помощью уникального атрибута:

resp = client.update_by_query(
    index="my-index-000001",
    query={
        "term": {
            "user.id": "kimchy"
        }
    },
    max_docs=1,
)
print(resp)
response = client.update_by_query(
  index: 'my-index-000001',
  body: {
    query: {
      term: {
        'user.id' => 'kimchy'
      }
    },
    max_docs: 1
  }
)
puts response
const response = await client.updateByQuery({
  index: "my-index-000001",
  query: {
    term: {
      "user.id": "kimchy",
    },
  },
  max_docs: 1,
});
console.log(response);
POST my-index-000001/_update_by_query
{
  "query": {
    "term": {
      "user.id": "kimchy"
    }
  },
  "max_docs": 1
}

Обновление источника документа

Обновление по запросу поддерживает скрипты для обновления источника документа. Например, следующий запрос увеличивает поле count для всех документов со значением user.id равным kimchy в my-index-000001:

resp = client.update_by_query(
    index="my-index-000001",
    script={
        "source": "ctx._source.count++",
        "lang": "painless"
    },
    query={
        "term": {
            "user.id": "kimchy"
        }
    },
)
print(resp)
response = client.update_by_query(
  index: 'my-index-000001',
  body: {
    script: {
      source: 'ctx._source.count++',
      lang: 'painless'
    },
    query: {
      term: {
        'user.id' => 'kimchy'
      }
    }
  }
)
puts response
const response = await client.updateByQuery({
  index: "my-index-000001",
  script: {
    source: "ctx._source.count++",
    lang: "painless",
  },
  query: {
    term: {
      "user.id": "kimchy",
    },
  },
});
console.log(response);
POST my-index-000001/_update_by_query
{
  "script": {
    "source": "ctx._source.count++",
    "lang": "painless"
  },
  "query": {
    "term": {
      "user.id": "kimchy"
    }
  }
}

Обратите внимание, что conflicts=proceed не указано в этом примере. В этом случае конфликт версий должен останавливать процесс, чтобы вы могли обработать ошибку.

Как и в API обновления, вы можете установить ctx.op для изменения выполняемой операции:

noop

Установите ctx.op = "noop", если ваш скрипт определит, что не нужно вносить никаких изменений. Операция обновления по запросу пропускает обновление документа и увеличивает счётчик noop.

delete

Установите ctx.op = "delete", если ваш скрипт определит, что документ следует удалить. Операция обновления по запросу удаляет документ и увеличивает счётчик deleted.

Обновление по запросу поддерживает только index, noop и delete. Установка ctx.op на другое значение является ошибкой. Установка любого другого поля в ctx является ошибкой. Этот API позволяет изменять только источник соответствующих документов, вы не можете их перемещать.

Обновление документов с помощью конвейера ingest

Обновление по запросу может использовать функцию конвейеров ingest, указав pipeline:

resp = client.ingest.put_pipeline(
    id="set-foo",
    description="sets foo",
    processors=[
        {
            "set": {
                "field": "foo",
                "value": "bar"
            }
        }
    ],
)
print(resp)

resp1 = client.update_by_query(
    index="my-index-000001",
    pipeline="set-foo",
)
print(resp1)
response = client.ingest.put_pipeline(
  id: 'set-foo',
  body: {
    description: 'sets foo',
    processors: [
      {
        set: {
          field: 'foo',
          value: 'bar'
        }
      }
    ]
  }
)
puts response

response = client.update_by_query(
  index: 'my-index-000001',
  pipeline: 'set-foo'
)
puts response
const response = await client.ingest.putPipeline({
  id: "set-foo",
  description: "sets foo",
  processors: [
    {
      set: {
        field: "foo",
        value: "bar",
      },
    },
  ],
});
console.log(response);

const response1 = await client.updateByQuery({
  index: "my-index-000001",
  pipeline: "set-foo",
});
console.log(response1);
PUT _ingest/pipeline/set-foo
{
  "description" : "sets foo",
  "processors" : [ {
      "set" : {
        "field": "foo",
        "value": "bar"
      }
  } ]
}
POST my-index-000001/_update_by_query?pipeline=set-foo
Получение статуса операций обновления по запросу

Вы можете получить статус всех выполняемых запросов обновления по запросу с помощью API задач:

$response = $client->tasks()->list();
resp = client.tasks.list(
    detailed=True,
    actions="*byquery",
)
print(resp)
response = client.tasks.list(
  detailed: true,
  actions: '*byquery'
)
puts response
res, err := es.Tasks.List(
	es.Tasks.List.WithActions("*byquery"),
	es.Tasks.List.WithDetailed(true),
)
fmt.Println(res, err)
const response = await client.tasks.list({
  detailed: "true",
  actions: "*byquery",
});
console.log(response);
GET _tasks?detailed=true&actions=*byquery

Ответы выглядят следующим образом:

{
  "nodes" : {
    "r1A2WoRbTwKZ516z6NEs5A" : {
      "name" : "r1A2WoR",
      "transport_address" : "127.0.0.1:9300",
      "host" : "127.0.0.1",
      "ip" : "127.0.0.1:9300",
      "attributes" : {
        "testattr" : "test",
        "portsfile" : "true"
      },
      "tasks" : {
        "r1A2WoRbTwKZ516z6NEs5A:36619" : {
          "node" : "r1A2WoRbTwKZ516z6NEs5A",
          "id" : 36619,
          "type" : "transport",
          "action" : "indices:data/write/update/byquery",
          "status" : {    
            "total" : 6154,
            "updated" : 3500,
            "created" : 0,
            "deleted" : 0,
            "batches" : 4,
            "version_conflicts" : 0,
            "noops" : 0,
            "retries": {
              "bulk": 0,
              "search": 0
            },
            "throttled_millis": 0
          },
          "description" : ""
        }
      }
    }
  }
}

Этот объект содержит фактический статус. Он похож на JSON-ответ с важным дополнением поля total. total — это общее количество операций, которые reindex ожидает выполнить. Вы можете оценить прогресс, сложив поля updated, created и deleted. Запрос завершится, когда их сумма будет равна значению поля total.

По идентификатору задачи вы можете получить информацию о задаче напрямую. Следующий пример извлекает информацию о задаче r1A2WoRbTwKZ516z6NEs5A:36619:

$params = [
    'task_id' => 'r1A2WoRbTwKZ516z6NEs5A:36619',
];
$response = $client->tasks()->get($params);
resp = client.tasks.get(
    task_id="r1A2WoRbTwKZ516z6NEs5A:36619",
)
print(resp)
response = client.tasks.get(
  task_id: 'r1A2WoRbTwKZ516z6NEs5A:36619'
)
puts response
res, err := es.Tasks.Get(
	"r1A2WoRbTwKZ516z6NEs5A:36619",
)
fmt.Println(res, err)
const response = await client.tasks.get({
  task_id: "r1A2WoRbTwKZ516z6NEs5A:36619",
});
console.log(response);
GET /_tasks/r1A2WoRbTwKZ516z6NEs5A:36619

Преимущества этого API заключаются в его интеграции с wait_for_completion=false для прозрачного возвращения статуса завершенных задач. Если задача завершена и для неё было установлено wait_for_completion=false, то она вернёт поле results или error. Стоимость этой функции — документ, который wait_for_completion=false создаёт в .tasks/task/${taskId}. Вам необходимо удалить этот документ.

Отмена операции обновления по запросу

Любое обновление по запросу может быть отменено с помощью API отмены задачи:

$params = [
    'task_id' => 'r1A2WoRbTwKZ516z6NEs5A:36619',
];
$response = $client->tasks()->cancel($params);
resp = client.tasks.cancel(
    task_id="r1A2WoRbTwKZ516z6NEs5A:36619",
)
print(resp)
response = client.tasks.cancel(
  task_id: 'r1A2WoRbTwKZ516z6NEs5A:36619'
)
puts response
res, err := es.Tasks.Cancel(
	es.Tasks.Cancel.WithTaskID("r1A2WoRbTwKZ516z6NEs5A:36619"),
)
fmt.Println(res, err)
const response = await client.tasks.cancel({
  task_id: "r1A2WoRbTwKZ516z6NEs5A:36619",
});
console.log(response);
POST _tasks/r1A2WoRbTwKZ516z6NEs5A:36619/_cancel

Идентификатор задачи можно найти с помощью API задач.

Отмена должна происходить быстро, но может занять несколько секунд. API статуса задачи выше продолжит отображать задачу обновления по запросу до тех пор, пока эта задача не проверит, что она отменена, и не завершит себя.

Изменение ограничения скорости запроса

Значение requests_per_second можно изменить для выполняемого обновления по запросу с помощью API _rethrottle:

$params = [
    'task_id' => 'r1A2WoRbTwKZ516z6NEs5A:36619',
];
$response = $client->updateByQueryRethrottle($params);
resp = client.update_by_query_rethrottle(
    task_id="r1A2WoRbTwKZ516z6NEs5A:36619",
    requests_per_second="-1",
)
print(resp)
response = client.update_by_query_rethrottle(
  task_id: 'r1A2WoRbTwKZ516z6NEs5A:36619',
  requests_per_second: -1
)
puts response
res, err := es.UpdateByQueryRethrottle(
	"r1A2WoRbTwKZ516z6NEs5A:36619",
	esapi.IntPtr(-1),
)
fmt.Println(res, err)
const response = await client.updateByQueryRethrottle({
  task_id: "r1A2WoRbTwKZ516z6NEs5A:36619",
  requests_per_second: "-1",
});
console.log(response);
POST _update_by_query/r1A2WoRbTwKZ516z6NEs5A:36619/_rethrottle?requests_per_second=-1

Идентификатор задачи можно найти с помощью API задач.

Как и при установке в API _update_by_query, requests_per_second может быть установлено в значение -1 для отключения ограничения скорости или в любое десятичное число, например, 1.7 или 12, для ограничения на этот уровень. Переограничение, которое ускоряет запрос, вступает в силу немедленно, но переограничение, которое замедляет запрос, вступит в силу после завершения текущей партии. Это предотвращает таймауты при прокрутке.

Ручное разделение на части

Разбейте обновление по запросу на части вручную, указав идентификатор части и общее количество частей в каждом запросе:

resp = client.update_by_query(
    index="my-index-000001",
    slice={
        "id": 0,
        "max": 2
    },
    script={
        "source": "ctx._source['extra'] = 'test'"
    },
)
print(resp)

resp1 = client.update_by_query(
    index="my-index-000001",
    slice={
        "id": 1,
        "max": 2
    },
    script={
        "source": "ctx._source['extra'] = 'test'"
    },
)
print(resp1)
response = client.update_by_query(
  index: 'my-index-000001',
  body: {
    slice: {
      id: 0,
      max: 2
    },
    script: {
      source: "ctx._source['extra'] = 'test'"
    }
  }
)
puts response

response = client.update_by_query(
  index: 'my-index-000001',
  body: {
    slice: {
      id: 1,
      max: 2
    },
    script: {
      source: "ctx._source['extra'] = 'test'"
    }
  }
)
puts response
const response = await client.updateByQuery({
  index: "my-index-000001",
  slice: {
    id: 0,
    max: 2,
  },
  script: {
    source: "ctx._source['extra'] = 'test'",
  },
});
console.log(response);

const response1 = await client.updateByQuery({
  index: "my-index-000001",
  slice: {
    id: 1,
    max: 2,
  },
  script: {
    source: "ctx._source['extra'] = 'test'",
  },
});
console.log(response1);
POST my-index-000001/_update_by_query
{
  "slice": {
    "id": 0,
    "max": 2
  },
  "script": {
    "source": "ctx._source['extra'] = 'test'"
  }
}
POST my-index-000001/_update_by_query
{
  "slice": {
    "id": 1,
    "max": 2
  },
  "script": {
    "source": "ctx._source['extra'] = 'test'"
  }
}

Что вы можете проверить, используя:

resp = client.indices.refresh()
print(resp)

resp1 = client.search(
    index="my-index-000001",
    size="0",
    q="extra:test",
    filter_path="hits.total",
)
print(resp1)
response = client.indices.refresh
puts response

response = client.search(
  index: 'my-index-000001',
  size: 0,
  q: 'extra:test',
  filter_path: 'hits.total'
)
puts response
const response = await client.indices.refresh();
console.log(response);

const response1 = await client.search({
  index: "my-index-000001",
  size: 0,
  q: "extra:test",
  filter_path: "hits.total",
});
console.log(response1);
GET _refresh
POST my-index-000001/_search?size=0&q=extra:test&filter_path=hits.total

Что приводит к разумному total, подобному этому:

{
  "hits": {
    "total": {
        "value": 120,
        "relation": "eq"
    }
  }
}
Использование автоматического разделения на части

Вы также можете позволить обновлению по запросу автоматически распараллеливать использование скроллинга частями для разделения на _id. Используйте slices для указания количества частей:

resp = client.update_by_query(
    index="my-index-000001",
    refresh=True,
    slices="5",
    script={
        "source": "ctx._source['extra'] = 'test'"
    },
)
print(resp)
response = client.update_by_query(
  index: 'my-index-000001',
  refresh: true,
  slices: 5,
  body: {
    script: {
      source: "ctx._source['extra'] = 'test'"
    }
  }
)
puts response
const response = await client.updateByQuery({
  index: "my-index-000001",
  refresh: "true",
  slices: 5,
  script: {
    source: "ctx._source['extra'] = 'test'",
  },
});
console.log(response);
POST my-index-000001/_update_by_query?refresh&slices=5
{
  "script": {
    "source": "ctx._source['extra'] = 'test'"
  }
}

Что также можно проверить с помощью:

resp = client.search(
    index="my-index-000001",
    size="0",
    q="extra:test",
    filter_path="hits.total",
)
print(resp)
response = client.search(
  index: 'my-index-000001',
  size: 0,
  q: 'extra:test',
  filter_path: 'hits.total'
)
puts response
const response = await client.search({
  index: "my-index-000001",
  size: 0,
  q: "extra:test",
  filter_path: "hits.total",
});
console.log(response);
POST my-index-000001/_search?size=0&q=extra:test&filter_path=hits.total

Что приводит к разумному total, например такому:

{
  "hits": {
    "total": {
        "value": 120,
        "relation": "eq"
    }
  }
}

Установка slices в значение auto позволит Elasticsearch выбрать количество частей. Этот параметр будет использовать одну часть на фрагмент до определённого предела. Если есть несколько потоков или индексов исходных данных, он выберет количество частей на основе индекса или базового индекса с наименьшим количеством фрагментов.

Добавление slices в _update_by_query просто автоматизирует ручной процесс, описанный выше, создавая подзапросы, что означает, что у него есть некоторые особенности:

  • Эти запросы можно увидеть в API задач. Эти подзапросы являются дочерними задачами задачи для запроса с slices.
  • Получение статуса задачи для запроса с slices содержит только статус завершённых частей.
  • Эти подзапросы индивидуально адресуемы для таких операций, как отмена и переограничение скорости.
  • Переограничение запроса с slices пропорционально переограничит незавершенный подзапрос.
  • Отмена запроса с slices отменяет каждый подзапрос.
  • Из-за особенностей slices каждый подзапрос не получит идеально равной части документов. Все документы будут обработаны, но некоторые части могут быть больше других. Ожидайте, что у более крупных частей будет более равномерное распределение.
  • Параметры, такие как requests_per_second и max_docs, в запросе с slices распределяются пропорционально каждому подзапросу. Объедините это с вышеупомянутым неравномерным распределением, и вы должны сделать вывод, что использование max_docs с slices, возможно, не приведёт к обновлению ровно max_docs документов.
  • Каждый подзапрос получает немного другой снимок исходных данных потока или индекса, хотя все они сделаны примерно в одно и то же время.
Добавление нового свойства

Предположим, вы создали индекс без динамической схемы, заполнили его данными, а затем добавили значение в схему, чтобы получить больше полей из данных:

$params = [
    'index' => 'test',
    'body' => [
        'mappings' => [
            'dynamic' => false,
            'properties' => [
                'text' => [
                    'type' => 'text',
                ],
            ],
        ],
    ],
];
$response = $client->indices()->create($params);
$params = [
    'index' => 'test',
    'body' => [
        'text' => 'words words',
        'flag' => 'bar',
    ],
];
$response = $client->index($params);
$params = [
    'index' => 'test',
    'body' => [
        'text' => 'words words',
        'flag' => 'foo',
    ],
];
$response = $client->index($params);
$params = [
    'index' => 'test',
    'body' => [
        'properties' => [
            'text' => [
                'type' => 'text',
            ],
            'flag' => [
                'type' => 'text',
                'analyzer' => 'keyword',
            ],
        ],
    ],
];
$response = $client->indices()->putMapping($params);
resp = client.indices.create(
    index="test",
    mappings={
        "dynamic": False,
        "properties": {
            "text": {
                "type": "text"
            }
        }
    },
)
print(resp)

resp1 = client.index(
    index="test",
    refresh=True,
    document={
        "text": "words words",
        "flag": "bar"
    },
)
print(resp1)

resp2 = client.index(
    index="test",
    refresh=True,
    document={
        "text": "words words",
        "flag": "foo"
    },
)
print(resp2)

resp3 = client.indices.put_mapping(
    index="test",
    properties={
        "text": {
            "type": "text"
        },
        "flag": {
            "type": "text",
            "analyzer": "keyword"
        }
    },
)
print(resp3)
response = client.indices.create(
  index: 'test',
  body: {
    mappings: {
      dynamic: false,
      properties: {
        text: {
          type: 'text'
        }
      }
    }
  }
)
puts response

response = client.index(
  index: 'test',
  refresh: true,
  body: {
    text: 'words words',
    flag: 'bar'
  }
)
puts response

response = client.index(
  index: 'test',
  refresh: true,
  body: {
    text: 'words words',
    flag: 'foo'
  }
)
puts response

response = client.indices.put_mapping(
  index: 'test',
  body: {
    properties: {
      text: {
        type: 'text'
      },
      flag: {
        type: 'text',
        analyzer: 'keyword'
      }
    }
  }
)
puts response
{
	res, err := es.Indices.Create(
		"test",
		es.Indices.Create.WithBody(strings.NewReader(`{
	  "mappings": {
	    "dynamic": false,
	    "properties": {
	      "text": {
	        "type": "text"
	      }
	    }
	  }
	}`)),
	)
	fmt.Println(res, err)
}

{
	res, err := es.Index(
		"test",
		strings.NewReader(`{
	  "text": "words words",
	  "flag": "bar"
	}`),
		es.Index.WithRefresh("true"),
		es.Index.WithPretty(),
	)
	fmt.Println(res, err)
}

{
	res, err := es.Index(
		"test",
		strings.NewReader(`{
	  "text": "words words",
	  "flag": "foo"
	}`),
		es.Index.WithRefresh("true"),
		es.Index.WithPretty(),
	)
	fmt.Println(res, err)
}

{
	res, err := es.Indices.PutMapping(
		[]string{"test"},
		strings.NewReader(`{
	  "properties": {
	    "text": {
	      "type": "text"
	    },
	    "flag": {
	      "type": "text",
	      "analyzer": "keyword"
	    }
	  }
	}`),
	)
	fmt.Println(res, err)
}
const response = await client.indices.create({
  index: "test",
  mappings: {
    dynamic: false,
    properties: {
      text: {
        type: "text",
      },
    },
  },
});
console.log(response);

const response1 = await client.index({
  index: "test",
  refresh: "true",
  document: {
    text: "words words",
    flag: "bar",
  },
});
console.log(response1);

const response2 = await client.index({
  index: "test",
  refresh: "true",
  document: {
    text: "words words",
    flag: "foo",
  },
});
console.log(response2);

const response3 = await client.indices.putMapping({
  index: "test",
  properties: {
    text: {
      type: "text",
    },
    flag: {
      type: "text",
      analyzer: "keyword",
    },
  },
});
console.log(response3);
PUT test
{
  "mappings": {
    "dynamic": false,   
    "properties": {
      "text": {"type": "text"}
    }
  }
}

POST test/_doc?refresh
{
  "text": "words words",
  "flag": "bar"
}
POST test/_doc?refresh
{
  "text": "words words",
  "flag": "foo"
}
PUT test/_mapping   
{
  "properties": {
    "text": {"type": "text"},
    "flag": {"type": "text", "analyzer": "keyword"}
  }
}

Это означает, что новые поля не будут индексированы, а только хранятся в _source.

Это обновляет схему, добавляя новое поле flag. Чтобы получить новое поле, необходимо повторно индексировать все документы с ним.

Поиск данных ничего не найдёт:

$params = [
    'index' => 'test',
    'body' => [
        'query' => [
            'match' => [
                'flag' => 'foo',
            ],
        ],
    ],
];
$response = $client->search($params);
resp = client.search(
    index="test",
    filter_path="hits.total",
    query={
        "match": {
            "flag": "foo"
        }
    },
)
print(resp)
response = client.search(
  index: 'test',
  filter_path: 'hits.total',
  body: {
    query: {
      match: {
        flag: 'foo'
      }
    }
  }
)
puts response
res, err := es.Search(
	es.Search.WithIndex("test"),
	es.Search.WithBody(strings.NewReader(`{
	  "query": {
	    "match": {
	      "flag": "foo"
	    }
	  }
	}`)),
	es.Search.WithFilterPath("hits.total"),
	es.Search.WithPretty(),
)
fmt.Println(res, err)
const response = await client.search({
  index: "test",
  filter_path: "hits.total",
  query: {
    match: {
      flag: "foo",
    },
  },
});
console.log(response);
POST test/_search?filter_path=hits.total
{
  "query": {
    "match": {
      "flag": "foo"
    }
  }
}
{
  "hits" : {
    "total": {
        "value": 0,
        "relation": "eq"
    }
  }
}

Но вы можете отправить запрос _update_by_query для получения нового отображения:

$params = [
    'index' => 'test',
];
$response = $client->updateByQuery($params);
$params = [
    'index' => 'test',
    'body' => [
        'query' => [
            'match' => [
                'flag' => 'foo',
            ],
        ],
    ],
];
$response = $client->search($params);
resp = client.update_by_query(
    index="test",
    refresh=True,
    conflicts="proceed",
)
print(resp)

resp1 = client.search(
    index="test",
    filter_path="hits.total",
    query={
        "match": {
            "flag": "foo"
        }
    },
)
print(resp1)
response = client.update_by_query(
  index: 'test',
  refresh: true,
  conflicts: 'proceed'
)
puts response

response = client.search(
  index: 'test',
  filter_path: 'hits.total',
  body: {
    query: {
      match: {
        flag: 'foo'
      }
    }
  }
)
puts response
{
	res, err := es.UpdateByQuery(
		[]string{"test"},
		es.UpdateByQuery.WithConflicts("proceed"),
		es.UpdateByQuery.WithRefresh(true),
	)
	fmt.Println(res, err)
}

{
	res, err := es.Search(
		es.Search.WithIndex("test"),
		es.Search.WithBody(strings.NewReader(`{
	  "query": {
	    "match": {
	      "flag": "foo"
	    }
	  }
	}`)),
		es.Search.WithFilterPath("hits.total"),
		es.Search.WithPretty(),
	)
	fmt.Println(res, err)
}
const response = await client.updateByQuery({
  index: "test",
  refresh: "true",
  conflicts: "proceed",
});
console.log(response);

const response1 = await client.search({
  index: "test",
  filter_path: "hits.total",
  query: {
    match: {
      flag: "foo",
    },
  },
});
console.log(response1);
POST test/_update_by_query?refresh&conflicts=proceed
POST test/_search?filter_path=hits.total
{
  "query": {
    "match": {
      "flag": "foo"
    }
  }
}
{
  "hits" : {
    "total": {
        "value": 1,
        "relation": "eq"
    }
  }
}

Вы можете сделать то же самое при добавлении поля в многопольное поле.

© 2023-2025 Elasticsearch
As of September 2024, Elasticsearch is available under a choice of three licenses: the Server Side Public License (SSPL), the Elastic License, or the AGPLv3 (OSI approved).
Elasticsearch and the Elasticsearch logo are trademarks of Elasticsearch B.V., registered in the U.S. and in other countries.
https://www.elastic.co/guide/en/elasticsearch/reference/8.17/docs-update-by-query.html

Spec-Zone.ru

Настройки Оффлайн Что нового Помощь О нас
Spec-Zone .ru
спецификации, руководства, описания, API