Java API
Устарело в 7.0.0.
TransportClient устарело в пользу Java High Level REST Client и будет удалено в Elasticsearch 8.0. Руководство по миграции описывает все шаги, необходимые для миграции.
X-Pack предоставляет Java-клиент под названием WatcherClient, который добавляет нативную поддержку Java для Watcher.
Для получения экземпляра WatcherClient, убедитесь, что вы сначала настроили XPackClient.
Установка XPackClient
Сначала вам нужно убедиться, что JAR-файл x-pack-transport-7.17.28 находится в классе. Вы можете извлечь этот JAR-файл из загруженного пакета X-Pack.
Если вы используете Maven для управления зависимостями, добавьте следующее в pom.xml:
<project ...>
<repositories>
<!-- add the elasticsearch repo -->
<repository>
<id>elasticsearch-releases</id>
<url>https://artifacts.elastic.co/maven</url>
<releases>
<enabled>true</enabled>
</releases>
<snapshots>
<enabled>false</enabled>
</snapshots>
</repository>
...
</repositories>
...
<dependencies>
<!-- add the x-pack jar as a dependency -->
<dependency>
<groupId>org.elasticsearch.client</groupId>
<artifactId>x-pack-transport</artifactId>
<version>7.17.28</version>
</dependency>
...
</dependencies>
...
</project> Если вы используете Gradle, добавьте зависимости в build.gradle:
repositories {
/* ... Any other repositories ... */
// Add the Elasticsearch Maven Repository
maven {
name "elastic"
url "https://artifacts.elastic.co/maven"
}
}
dependencies {
// Provide the x-pack jar on the classpath for compilation and at runtime
compile "org.elasticsearch.client:x-pack-transport:7.17.28"
/* ... */
} Вы также можете загрузить JAR-файл X-Pack Transport вручную, непосредственно из нашего репозитория Maven.
Получение WatcherClient
Чтобы получить экземпляр WatcherClient, сначала нужно создать XPackClient. XPackClient — это обертка вокруг стандартного Java Elasticsearch Client:
import org.elasticsearch.client.transport.TransportClient;
import org.elasticsearch.xpack.client.PreBuiltXPackTransportClient;
import org.elasticsearch.xpack.core.XPackClient;
import org.elasticsearch.xpack.core.XPackPlugin;
import org.elasticsearch.core.watcher.client.WatcherClient;
...
TransportClient client = new PreBuiltXPackTransportClient(Settings.builder()
.put("cluster.name", "myClusterName")
...
.build())
.addTransportAddress(new TransportAddress(InetAddress.getByName("localhost"), 9300));
XPackClient xpackClient = new XPackClient(client);
WatcherClient watcherClient = xpackClient.watcher(); Создание или обновление API watch
API создания или обновления watch либо регистрирует новое watch в Watcher, либо обновляет существующее. После регистрации новый документ будет добавлен в индекс .watches, представляющий watch, и триггер watch будет немедленно зарегистрирован в соответствующем движке триггеров (обычно планировщике, для триггера schedule).
Добавление watch должно выполняться только через этот API. Не добавляйте watch напрямую в индекс .watches с помощью API индекса Elasticsearch. Если включены функции безопасности Elasticsearch, убедитесь, что никому не предоставлены привилегии write для индекса .watches.
Следующий пример добавляет watch с идентификатором my-watch, который имеет следующие характеристики:
- Расписание watch запускается каждую минуту.
- Входной поиск watch ищет любые HTTP-ответы 404, которые произошли в течение последних пяти минут.
- Условие watch проверяет, были ли найдены какие-либо совпадения.
- При обнаружении совпадений watch отправляет администратору электронное письмо.
WatchSourceBuilder watchSourceBuilder = WatchSourceBuilders.watchBuilder();
// Set the trigger
watchSourceBuilder.trigger(TriggerBuilders.schedule(Schedules.cron("0 0/1 * * * ?")));
// Create the search request to use for the input
SearchRequest request = Requests.searchRequest("idx").source(searchSource()
.query(boolQuery()
.must(matchQuery("response", 404))
.filter(rangeQuery("date").gt("{{ctx.trigger.scheduled_time}}"))
.filter(rangeQuery("date").lt("{{ctx.execution_time}}"))
));
// Create the search input
SearchInput input = new SearchInput(new WatcherSearchTemplateRequest(new String[]{"idx"}, null, SearchType.DEFAULT,
WatcherSearchTemplateRequest.DEFAULT_INDICES_OPTIONS, new BytesArray(request.source().toString())), null, null, null);
// Set the input
watchSourceBuilder.input(input);
// Set the condition
watchSourceBuilder.condition(new ScriptCondition(new Script("ctx.payload.hits.total > 1")));
// Create the email template to use for the action
EmailTemplate.Builder emailBuilder = EmailTemplate.builder();
emailBuilder.to("someone@domain.host.com");
emailBuilder.subject("404 recently encountered");
EmailAction.Builder emailActionBuilder = EmailAction.builder(emailBuilder.build());
// Add the action
watchSourceBuilder.addAction("email_someone", emailActionBuilder);
PutWatchResponse putWatchResponse = watcherClient.preparePutWatch("my-watch")
.setSource(watchSourceBuilder)
.get(); Хотя приведенный выше фрагмент демонстрирует все конкретные классы, которые составляют наше watch, использование доступных классов-билдеров вместе со статическими импортами может значительно упростить и уплотнить ваш код:
PutWatchResponse putWatchResponse2 = watcherClient.preparePutWatch("my-watch")
.setSource(watchBuilder()
.trigger(schedule(cron("0 0/1 * * * ?")))
.input(searchInput(new WatcherSearchTemplateRequest(new String[]{"idx"}, null, SearchType.DEFAULT,
WatcherSearchTemplateRequest.DEFAULT_INDICES_OPTIONS, searchSource()
.query(boolQuery()
.must(matchQuery("response", 404))
.filter(rangeQuery("date").gt("{{ctx.trigger.scheduled_time}}"))
.filter(rangeQuery("date").lt("{{ctx.execution_time}}"))
).buildAsBytes())))
.condition(compareCondition("ctx.payload.hits.total", CompareCondition.Op.GT, 1L))
.addAction("email_someone", emailAction(EmailTemplate.builder()
.to("someone@domain.host.com")
.subject("404 recently encountered"))))
.get(); - Используйте классы
TriggerBuildersиSchedulesдля определения триггера - Используйте класс
InputBuildersдля определения входных данных - Используйте класс
ConditionBuildersдля определения условия - Используйте
ActionBuildersдля определения действий
API получения watch
Этот API получает watch по его идентификатору.
Следующий пример получает watch с идентификатором my-watch:
GetWatchResponse getWatchResponse = watcherClient.prepareGetWatch("my-watch").get(); Вы можете получить определение watch, обратившись к источнику ответа:
XContentSource source = getWatchResponse.getSource();
XContentSource предоставляет вам методы для изучения источника:
Map<String, Object> map = source.getAsMap();
Или получить определенное значение, связанное с известным ключом:
String host = source.getValue("input.http.request.host"); API удаления watch
API удаления watch удаляет watch (идентифицированный по его id) из Watcher. После удаления документ, представляющий watch в индексе .watches, исчезает, и он больше никогда не будет выполняться.
Обратите внимание, что удаление watch не удаляет записи о выполнении watch, связанные с этим watch, из истории watch.
Удаление watch должно выполняться только через этот API. Не удаляйте watch напрямую из индекса .watches с помощью API удаления документа Elasticsearch. Если включены функции безопасности Elasticsearch, убедитесь, что никому не предоставлены привилегии write для индекса .watches.
Следующий пример удаляет watch с идентификатором my-watch:
DeleteWatchResponse deleteWatchResponse = watcherClient.prepareDeleteWatch("my-watch").get(); API выполнения watch
Этот API позволяет выполнять watch, хранящийся в индексе .watches, по запросу. Он может использоваться для тестирования watch без выполнения всех его действий или игнорирования его условия. Ответ содержит BytesReference, который представляет запись, которая будет записана в индекс .watcher-history.
Следующий пример выполняет watch с именем my-watch
ExecuteWatchResponse executeWatchResponse = watcherClient.prepareExecuteWatch("my-watch")
// execute the actions, ignoring the watch condition
.setIgnoreCondition(true)
// A map containing alternative input to use instead of the output of
// the watch's input
.setAlternativeInput(new HashMap<String, Object>())
// Trigger data to use (Note that "scheduled_time" is not provided to the
// ctx.trigger by this execution method so you may want to include it here)
.setTriggerData(new HashMap<String, Object>())
// Simulating the "email_admin" action while ignoring its throttle state. Use
// "_all" to set the action execution mode to all actions
.setActionMode("_all", ActionExecutionMode.FORCE_SIMULATE)
// If the execution of this watch should be written to the `.watcher-history`
// index and reflected in the persisted Watch
.setRecordExecution(false)
// Indicates whether the watch should execute in debug mode. In debug mode the
// returned watch record will hold the execution vars
.setDebug(true)
.get(); После возвращения ответа вы можете изучить его, получив исходные данные записи выполнения:
Класс XContentSource предоставляет удобные методы для изучения источника
XContentSource source = executeWatchResponse.getRecordSource();
String actionId = source.getValue("result.actions.0.id"); API подтверждения watch
Подтверждение watch позволяет вручную ограничивать выполнение действий watch. Состояние подтверждения действия хранится в структуре status.actions.<id>.ack.state.
Текущий статус watch и состояние его действий возвращаются в качестве части ответа API получения watch:
GetWatchResponse getWatchResponse = watcherClient.prepareGetWatch("my-watch").get();
State state = getWatchResponse.getStatus().actionStatus("my-action").ackStatus().state(); Состояние действия вновь созданного watch равно awaits_successful_execution. Когда watch выполняется и его условие выполняется, состояние изменяется на ackable. Подтверждение действия устанавливает состояние на acked.
Когда состояние действия установлено на acked, дальнейшие выполнения этого действия ограничены до тех пор, пока его состояние не будет сброшено до awaits_successful_execution. Это происходит, когда условие watch больше не выполняется (условие оценивается как false).
Следующий фрагмент показывает, как подтвердить действие. Вы указываете идентификаторы watch и действия, которые вы хотите подтвердить — в этом примере my-watch и my-action:
AckWatchResponse ackResponse = watcherClient.prepareAckWatch("my-watch").setActionIds("my-action").get(); В ответ на этот запрос возвращаются статус watch и состояние действия, которые можно получить из объекта AckWatchResponse:
WatchStatus status = ackResponse.getStatus();
ActionStatus actionStatus = status.actionStatus("my-action");
ActionStatus.AckStatus ackStatus = actionStatus.ackStatus();
ActionStatus.AckStatus.State ackState = ackStatus.state(); Вы можете подтвердить несколько действий:
AckWatchResponse ackResponse = watcherClient.prepareAckWatch("my-watch")
.setActionIds("action1", "action2")
.get(); Чтобы подтвердить все действия watch, укажите только идентификатор watch:
AckWatchResponse ackResponse = watcherClient.prepareAckWatch("my-watch").get(); API активации watch
Watch может быть либо активным, либо неактивным. Этот API позволяет активировать в настоящее время неактивный watch.
Статус неактивного watch возвращается вместе с определением watch при вызове API получения watch:
GetWatchResponse getWatchResponse = watcherClient.prepareGetWatch("my-watch").get();
boolean active = getWatchResponse.getStatus().state().isActive(); Следующий фрагмент показывает, как вы можете активировать watch:
ActivateWatchResponse activateResponse = watcherClient.prepareActivateWatch("my-watch", true).get();
boolean active = activateResponse.getStatus().state().isActive(); Новый статус watch возвращается как часть его общего статуса.
API деактивации watch
Watch может быть либо активным, либо неактивным. Этот API позволяет деактивировать в настоящее время активный watch.
Статус активного watch возвращается вместе с определением watch при вызове API получения watch:
GetWatchResponse getWatchResponse = watcherClient.prepareGetWatch("my-watch").get();
boolean active = getWatchResponse.getStatus().state().isActive(); Следующий фрагмент показывает, как вы можете деактивировать watch:
ActivateWatchResponse activateResponse = watcherClient.prepareActivateWatch("my-watch", false).get();
boolean active = activateResponse.getStatus().state().isActive(); Новый статус watch возвращается как часть его общего статуса.
API статистики Watcher
API stats возвращает текущие метрики Watcher. Вы можете контролировать, какие метрики возвращает этот API, используя параметр metric.
Следующий пример обращается к API stats:
WatcherStatsResponse watcherStatsResponse = watcherClient.prepareWatcherStats().get();
Успешный вызов возвращает структуру ответа, к которой можно обратиться следующим образом:
WatcherBuild build = watcherStatsResponse.getBuild();
// The current size of the watcher execution queue
long executionQueueSize = watcherStatsResponse.getThreadPoolQueueSize();
// The maximum size the watch execution queue has grown to
long executionQueueMaxSize = watcherStatsResponse.getThreadPoolQueueSize();
// The total number of watches registered in the system
long totalNumberOfWatches = watcherStatsResponse.getWatchesCount();
// {watcher} state (STARTING,STOPPED or STARTED)
WatcherState watcherState = watcherStatsResponse.getWatcherState(); API службы Watcher
API Watcher service позволяет управлять жизненным циклом службы Watcher. Следующий пример запускает службу watcher:
WatcherServiceResponse watcherServiceResponse = watcherClient.prepareWatchService().start().get();
Следующий пример останавливает службу watcher:
WatcherServiceResponse watcherServiceResponse = watcherClient.prepareWatchService().stop().get();
Следующий пример перезапускает службу watcher:
WatcherServiceResponse watcherServiceResponse = watcherClient.prepareWatchService().restart().get();
© 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/7.17/api-java.html