1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16 package org.codelibs.fess.es.client;
17
18 import static org.codelibs.core.stream.StreamUtil.split;
19 import static org.codelibs.core.stream.StreamUtil.stream;
20 import static org.codelibs.fesen.runner.FesenRunner.newConfigs;
21
22 import java.io.File;
23 import java.io.IOException;
24 import java.net.InetAddress;
25 import java.net.UnknownHostException;
26 import java.nio.charset.StandardCharsets;
27 import java.text.SimpleDateFormat;
28 import java.util.ArrayList;
29 import java.util.Arrays;
30 import java.util.Collections;
31 import java.util.Date;
32 import java.util.HashMap;
33 import java.util.List;
34 import java.util.Map;
35 import java.util.Map.Entry;
36 import java.util.function.BiConsumer;
37 import java.util.function.BiFunction;
38 import java.util.function.Function;
39 import java.util.regex.Pattern;
40 import java.util.stream.Collectors;
41
42 import javax.annotation.PostConstruct;
43 import javax.annotation.PreDestroy;
44
45 import org.apache.logging.log4j.LogManager;
46 import org.apache.logging.log4j.Logger;
47 import org.codelibs.core.beans.util.BeanUtil;
48 import org.codelibs.core.exception.ResourceNotFoundRuntimeException;
49 import org.codelibs.core.io.FileUtil;
50 import org.codelibs.core.io.ResourceUtil;
51 import org.codelibs.core.lang.StringUtil;
52 import org.codelibs.core.lang.ThreadUtil;
53 import org.codelibs.curl.CurlResponse;
54 import org.codelibs.fesen.FesenException;
55 import org.codelibs.fesen.FesenStatusException;
56 import org.codelibs.fesen.action.ActionFuture;
57 import org.codelibs.fesen.action.ActionListener;
58 import org.codelibs.fesen.action.ActionRequest;
59 import org.codelibs.fesen.action.ActionResponse;
60 import org.codelibs.fesen.action.ActionType;
61 import org.codelibs.fesen.action.DocWriteRequest;
62 import org.codelibs.fesen.action.DocWriteRequest.OpType;
63 import org.codelibs.fesen.action.DocWriteResponse.Result;
64 import org.codelibs.fesen.action.admin.cluster.health.ClusterHealthResponse;
65 import org.codelibs.fesen.action.admin.indices.alias.IndicesAliasesRequestBuilder;
66 import org.codelibs.fesen.action.admin.indices.create.CreateIndexResponse;
67 import org.codelibs.fesen.action.admin.indices.exists.indices.IndicesExistsResponse;
68 import org.codelibs.fesen.action.admin.indices.flush.FlushResponse;
69 import org.codelibs.fesen.action.admin.indices.get.GetIndexResponse;
70 import org.codelibs.fesen.action.admin.indices.mapping.get.GetMappingsResponse;
71 import org.codelibs.fesen.action.admin.indices.refresh.RefreshResponse;
72 import org.codelibs.fesen.action.bulk.BulkItemResponse;
73 import org.codelibs.fesen.action.bulk.BulkItemResponse.Failure;
74 import org.codelibs.fesen.action.bulk.BulkRequest;
75 import org.codelibs.fesen.action.bulk.BulkRequestBuilder;
76 import org.codelibs.fesen.action.bulk.BulkResponse;
77 import org.codelibs.fesen.action.delete.DeleteRequest;
78 import org.codelibs.fesen.action.delete.DeleteRequestBuilder;
79 import org.codelibs.fesen.action.delete.DeleteResponse;
80 import org.codelibs.fesen.action.explain.ExplainRequest;
81 import org.codelibs.fesen.action.explain.ExplainRequestBuilder;
82 import org.codelibs.fesen.action.explain.ExplainResponse;
83 import org.codelibs.fesen.action.fieldcaps.FieldCapabilitiesRequest;
84 import org.codelibs.fesen.action.fieldcaps.FieldCapabilitiesRequestBuilder;
85 import org.codelibs.fesen.action.fieldcaps.FieldCapabilitiesResponse;
86 import org.codelibs.fesen.action.get.GetRequest;
87 import org.codelibs.fesen.action.get.GetRequestBuilder;
88 import org.codelibs.fesen.action.get.GetResponse;
89 import org.codelibs.fesen.action.get.MultiGetRequest;
90 import org.codelibs.fesen.action.get.MultiGetRequestBuilder;
91 import org.codelibs.fesen.action.get.MultiGetResponse;
92 import org.codelibs.fesen.action.index.IndexRequest;
93 import org.codelibs.fesen.action.index.IndexRequestBuilder;
94 import org.codelibs.fesen.action.index.IndexResponse;
95 import org.codelibs.fesen.action.search.ClearScrollRequest;
96 import org.codelibs.fesen.action.search.ClearScrollRequestBuilder;
97 import org.codelibs.fesen.action.search.ClearScrollResponse;
98 import org.codelibs.fesen.action.search.MultiSearchRequest;
99 import org.codelibs.fesen.action.search.MultiSearchRequestBuilder;
100 import org.codelibs.fesen.action.search.MultiSearchResponse;
101 import org.codelibs.fesen.action.search.SearchPhaseExecutionException;
102 import org.codelibs.fesen.action.search.SearchRequest;
103 import org.codelibs.fesen.action.search.SearchRequestBuilder;
104 import org.codelibs.fesen.action.search.SearchResponse;
105 import org.codelibs.fesen.action.search.SearchScrollRequest;
106 import org.codelibs.fesen.action.search.SearchScrollRequestBuilder;
107 import org.codelibs.fesen.action.support.WriteRequest.RefreshPolicy;
108 import org.codelibs.fesen.action.support.master.AcknowledgedResponse;
109 import org.codelibs.fesen.action.termvectors.MultiTermVectorsRequest;
110 import org.codelibs.fesen.action.termvectors.MultiTermVectorsRequestBuilder;
111 import org.codelibs.fesen.action.termvectors.MultiTermVectorsResponse;
112 import org.codelibs.fesen.action.termvectors.TermVectorsRequest;
113 import org.codelibs.fesen.action.termvectors.TermVectorsRequestBuilder;
114 import org.codelibs.fesen.action.termvectors.TermVectorsResponse;
115 import org.codelibs.fesen.action.update.UpdateRequest;
116 import org.codelibs.fesen.action.update.UpdateRequestBuilder;
117 import org.codelibs.fesen.action.update.UpdateResponse;
118 import org.codelibs.fesen.client.AdminClient;
119 import org.codelibs.fesen.client.Client;
120 import org.codelibs.fesen.client.HttpClient;
121 import org.codelibs.fesen.cluster.metadata.MappingMetadata;
122 import org.codelibs.fesen.common.collect.ImmutableOpenMap;
123 import org.codelibs.fesen.common.document.DocumentField;
124 import org.codelibs.fesen.common.settings.Settings;
125 import org.codelibs.fesen.common.settings.Settings.Builder;
126 import org.codelibs.fesen.common.xcontent.XContentType;
127 import org.codelibs.fesen.core.TimeValue;
128 import org.codelibs.fesen.index.query.InnerHitBuilder;
129 import org.codelibs.fesen.index.query.QueryBuilder;
130 import org.codelibs.fesen.index.query.QueryBuilders;
131 import org.codelibs.fesen.rest.RestStatus;
132 import org.codelibs.fesen.runner.FesenRunner;
133 import org.codelibs.fesen.runner.FesenRunner.Configs;
134 import org.codelibs.fesen.search.SearchHit;
135 import org.codelibs.fesen.search.SearchHits;
136 import org.codelibs.fesen.search.aggregations.AggregationBuilders;
137 import org.codelibs.fesen.search.aggregations.bucket.filter.FilterAggregationBuilder;
138 import org.codelibs.fesen.search.aggregations.bucket.terms.TermsAggregationBuilder;
139 import org.codelibs.fesen.search.collapse.CollapseBuilder;
140 import org.codelibs.fesen.search.fetch.subphase.highlight.HighlightBuilder;
141 import org.codelibs.fesen.threadpool.ThreadPool;
142 import org.codelibs.fess.Constants;
143 import org.codelibs.fess.entity.FacetInfo;
144 import org.codelibs.fess.entity.GeoInfo;
145 import org.codelibs.fess.entity.HighlightInfo;
146 import org.codelibs.fess.entity.PingResponse;
147 import org.codelibs.fess.entity.QueryContext;
148 import org.codelibs.fess.entity.SearchRequestParams.SearchRequestType;
149 import org.codelibs.fess.exception.FessSystemException;
150 import org.codelibs.fess.exception.InvalidQueryException;
151 import org.codelibs.fess.exception.ResultOffsetExceededException;
152 import org.codelibs.fess.exception.SearchQueryException;
153 import org.codelibs.fess.helper.DocumentHelper;
154 import org.codelibs.fess.helper.QueryHelper;
155 import org.codelibs.fess.mylasta.direction.FessConfig;
156 import org.codelibs.fess.util.BooleanFunction;
157 import org.codelibs.fess.util.ComponentUtil;
158 import org.codelibs.fess.util.DocMap;
159 import org.dbflute.exception.IllegalBehaviorStateException;
160 import org.dbflute.optional.OptionalEntity;
161 import org.lastaflute.core.message.UserMessages;
162 import org.lastaflute.di.exception.ContainerInitFailureException;
163
164 import com.fasterxml.jackson.core.type.TypeReference;
165 import com.fasterxml.jackson.databind.ObjectMapper;
166 import com.google.common.io.BaseEncoding;
167
168 public class SearchEngineClient implements Client {
169
170 private static final Logger logger = LogManager.getLogger(SearchEngineClient.class);
171
172 private static final String LOG_INDEX_PREFIX = "fess_log";
173
174 private static final String USER_INDEX_PREFIX = ".fess_user";
175
176 private static final String CONFIG_INDEX_PREFIX = ".fess_config";
177
178 protected FesenRunner runner;
179
180 protected Client client;
181
182 protected Map<String, String> settings;
183
184 protected String indexConfigPath = "fess_indices";
185
186 protected List<String> indexConfigList = new ArrayList<>();
187
188 protected Map<String, List<String>> configListMap = new HashMap<>();
189
190 protected String scrollForSearch = "1m";
191
192 protected int sizeForDelete = 100;
193
194 protected String scrollForDelete = "1m";
195
196 protected int sizeForUpdate = 100;
197
198 protected String scrollForUpdate = "1m";
199
200 protected int maxConfigSyncStatusRetry = 10;
201
202 protected int maxEsStatusRetry = 60;
203
204 protected String clusterName = "fesen";
205
206 public void addIndexConfig(final String path) {
207 indexConfigList.add(path);
208 }
209
210 public void addConfigFile(final String index, final String path) {
211 configListMap.computeIfAbsent(index, k -> new ArrayList<>()).add(path);
212 }
213
214 public void setSettings(final Map<String, String> settings) {
215 this.settings = settings;
216 }
217
218 public String getStatus() {
219 return admin().cluster().prepareHealth().execute().actionGet(ComponentUtil.getFessConfig().getIndexHealthTimeout()).getStatus()
220 .name();
221 }
222
223 public void setRunner(final FesenRunner runner) {
224 this.runner = runner;
225 }
226
227 public boolean isEmbedded() {
228 return this.runner != null;
229 }
230
231 protected InetAddress getInetAddressByName(final String host) {
232 try {
233 return InetAddress.getByName(host);
234 } catch (final UnknownHostException e) {
235 throw new FessSystemException("Failed to resolve the hostname: " + host, e);
236 }
237 }
238
239 @PostConstruct
240 public void open() {
241 if (logger.isDebugEnabled()) {
242 logger.debug("Initialize {}", this.getClass().getSimpleName());
243 }
244 final FessConfig fessConfig = ComponentUtil.getFessConfig();
245
246 String httpAddress = System.getProperty(Constants.FESS_ES_HTTP_ADDRESS);
247 if (StringUtil.isBlank(httpAddress) && (runner == null)) {
248 switch (fessConfig.getFesenType()) {
249 case Constants.FESEN_TYPE_CLOUD:
250 case Constants.FESEN_TYPE_AWS:
251 httpAddress = org.codelibs.fess.util.ResourceUtil.getFesenHttpUrl();
252 break;
253 default:
254 runner = new FesenRunner();
255 final Configs config = newConfigs().clusterName(clusterName).numOfNode(1).useLogger();
256 final String esDir = System.getProperty("fess.es.dir");
257 if (esDir != null) {
258 config.basePath(esDir);
259 }
260 config.disableESLogger();
261 runner.onBuild((number, settingsBuilder) -> {
262 final File pluginDir = new File(esDir, "plugins");
263 if (pluginDir.isDirectory()) {
264 settingsBuilder.put("path.plugins", pluginDir.getAbsolutePath());
265 } else {
266 settingsBuilder.put("path.plugins", new File(System.getProperty("user.dir"), "plugins").getAbsolutePath());
267 }
268 if (settings != null) {
269 settingsBuilder.putProperties(settings, s -> s);
270 }
271 });
272 runner.build(config);
273
274 final int port = runner.node().settings().getAsInt("http.port", 9200);
275 httpAddress = "http://localhost:" + port;
276 logger.warn("Embedded Fesen is running. This configuration is not recommended for production use.");
277 break;
278 }
279 }
280 client = createHttpClient(fessConfig, httpAddress);
281
282 if (StringUtil.isNotBlank(httpAddress)) {
283 System.setProperty(Constants.FESS_ES_HTTP_ADDRESS, httpAddress);
284 }
285
286 waitForYellowStatus(fessConfig);
287
288 indexConfigList.forEach(configName -> {
289 final String[] values = configName.split("/");
290 if (values.length == 2) {
291 final String configIndex = values[0];
292 final String configType = values[1];
293
294 final boolean isFessIndex = "fess".equals(configIndex);
295 final String indexName;
296 if (isFessIndex) {
297 final boolean exists = existsIndex(fessConfig.getIndexDocumentUpdateIndex());
298 if (!exists) {
299 indexName = generateNewIndexName(configIndex);
300 createIndex(configIndex, indexName);
301 createAlias(configIndex, indexName);
302 } else {
303 client.admin().cluster().prepareHealth(fessConfig.getIndexDocumentUpdateIndex()).setWaitForYellowStatus().execute()
304 .actionGet(fessConfig.getIndexIndicesTimeout());
305 final GetIndexResponse response =
306 client.admin().indices().prepareGetIndex().addIndices(fessConfig.getIndexDocumentUpdateIndex()).execute()
307 .actionGet(fessConfig.getIndexIndicesTimeout());
308 final String[] indices = response.indices();
309 if (indices.length == 1) {
310 indexName = indices[0];
311 } else {
312 indexName = configIndex;
313 }
314 }
315 } else {
316 if (configIndex.startsWith(CONFIG_INDEX_PREFIX)) {
317 final String name = fessConfig.getIndexConfigIndex();
318 indexName = configIndex.replaceFirst(Pattern.quote(CONFIG_INDEX_PREFIX), name);
319 } else if (configIndex.startsWith(USER_INDEX_PREFIX)) {
320 final String name = fessConfig.getIndexUserIndex();
321 indexName = configIndex.replaceFirst(Pattern.quote(CONFIG_INDEX_PREFIX), name);
322 } else if (configIndex.startsWith(LOG_INDEX_PREFIX)) {
323 final String name = fessConfig.getIndexLogIndex();
324 indexName = configIndex.replaceFirst(Pattern.quote(CONFIG_INDEX_PREFIX), name);
325 } else {
326 throw new FessSystemException("Unknown config index: " + configIndex);
327 }
328 final boolean exists = existsIndex(indexName);
329 if (!exists) {
330 createIndex(configIndex, indexName);
331 createAlias(configIndex, indexName);
332 }
333 }
334
335 addMapping(configIndex, configType, indexName);
336 } else {
337 logger.warn("Invalid index config name: {}", configName);
338 }
339 });
340 }
341
342 protected Client createHttpClient(final FessConfig fessConfig, final String host) {
343 final Builder builder = Settings.builder().putList("http.hosts", host).put("processors", fessConfig.availableProcessors());
344 if (StringUtil.isNotBlank(fessConfig.getFesenUsername()) && StringUtil.isNotBlank(fessConfig.getFesenPassword())) {
345 builder.put(Constants.FESEN_USERNAME, fessConfig.getFesenUsername());
346 builder.put(Constants.FESEN_PASSWORD, fessConfig.getFesenPassword());
347 }
348 return new HttpClient(builder.build(), null);
349 }
350
351 public boolean existsIndex(final String indexName) {
352 final FessConfig fessConfig = ComponentUtil.getFessConfig();
353 boolean exists = false;
354 try {
355 final IndicesExistsResponse response =
356 client.admin().indices().prepareExists(indexName).execute().actionGet(fessConfig.getIndexSearchTimeout());
357 exists = response.isExists();
358 } catch (final Exception e) {
359
360 }
361 return exists;
362 }
363
364 public boolean reindex(final String fromIndex, final String toIndex, final boolean waitForCompletion) {
365 final FessConfig fessConfig = ComponentUtil.getFessConfig();
366 final String source = fessConfig.getIndexReindexBody()
367 .replace("__SOURCE_INDEX__", fromIndex)
368 .replace("__DEST_INDEX__", toIndex)
369 .replace("__SCRIPT_SOURCE__", ComponentUtil.getLanguageHelper().getReindexScriptSource());
370 final String refresh = StringUtil.isNotBlank(fessConfig.getIndexReindexRefresh()) ? fessConfig.getIndexReindexRefresh() : null;
371 final String requestsPerSecond =
372 StringUtil.isNotBlank(fessConfig.getIndexReindexRequestsPerSecond()) ? fessConfig.getIndexReindexRequestsPerSecond() : null;
373 final String scroll = StringUtil.isNotBlank(fessConfig.getIndexReindexScroll()) ? fessConfig.getIndexReindexScroll() : null;
374 final String maxDocs = StringUtil.isNotBlank(fessConfig.getIndexReindexMaxDocs()) ? fessConfig.getIndexReindexMaxDocs() : null;
375 try (CurlResponse response = ComponentUtil.getCurlHelper().post("/_reindex").param("refresh", refresh)
376 .param("requests_per_second", requestsPerSecond).param("scroll", scroll).param("max_docs", maxDocs)
377 .param("wait_for_completion", Boolean.toString(waitForCompletion)).body(source).execute()) {
378 if (response.getHttpStatusCode() == 200) {
379 return true;
380 }
381 logger.warn("Failed to reindex from {} to {}", fromIndex, toIndex);
382 } catch (final IOException e) {
383 logger.warn("Failed to reindex from {} to {}", fromIndex, toIndex, e);
384 }
385 return false;
386 }
387
388 public boolean createIndex(final String index, final String indexName) {
389 final FessConfig fessConfig = ComponentUtil.getFessConfig();
390 return createIndex(index, indexName, fessConfig.getIndexNumberOfShards(), fessConfig.getIndexAutoExpandReplicas(), true);
391 }
392
393 public boolean createIndex(final String index, final String indexName, final String numberOfShards, final String autoExpandReplicas,
394 final boolean uploadConfig) {
395 final FessConfig fessConfig = ComponentUtil.getFessConfig();
396
397 final String fesenType = fessConfig.getFesenType();
398 if (uploadConfig) {
399 switch (fesenType) {
400 case Constants.FESEN_TYPE_CLOUD:
401 case Constants.FESEN_TYPE_AWS:
402
403 break;
404 default:
405 waitForConfigSyncStatus();
406 sendConfigFiles(index);
407 break;
408 }
409 }
410
411 final String indexConfigFile = getResourcePath(indexConfigPath, fesenType, "/" + index + ".json");
412 try {
413 final String source = readIndexSetting(fesenType, indexConfigFile, numberOfShards, autoExpandReplicas);
414 final CreateIndexResponse indexResponse = client.admin().indices().prepareCreate(indexName).setSource(source, XContentType.JSON)
415 .execute().actionGet(fessConfig.getIndexIndicesTimeout());
416 if (indexResponse.isAcknowledged()) {
417 logger.info("Created {} index.", indexName);
418 return true;
419 }
420 if (logger.isDebugEnabled()) {
421 logger.debug("Failed to create {} index.", indexName);
422 }
423 } catch (final Exception e) {
424 logger.warn("{} is not found.", indexConfigFile, e);
425 }
426
427 return false;
428 }
429
430 protected String readIndexSetting(final String fesenType, final String indexConfigFile, final String numberOfShards,
431 final String autoExpandReplicas) {
432 final FessConfig fessConfig = ComponentUtil.getFessConfig();
433 String source = FileUtil.readUTF8(indexConfigFile);
434 String dictionaryPath = System.getProperty("fess.dictionary.path", StringUtil.EMPTY);
435 if (StringUtil.isNotBlank(dictionaryPath) && !dictionaryPath.endsWith("/")) {
436 dictionaryPath = dictionaryPath + "/";
437 }
438 source = source.replaceAll(Pattern.quote("${fess.dictionary.path}"), dictionaryPath);
439 source = source.replaceAll(Pattern.quote("${fess.index.codec}"), fessConfig.getIndexCodec());
440 source = source.replaceAll(Pattern.quote("${fess.index.number_of_shards}"), numberOfShards);
441 source = source.replaceAll(Pattern.quote("${fess.index.auto_expand_replicas}"), autoExpandReplicas);
442 return source;
443 }
444
445 protected String getResourcePath(final String basePath, final String type, final String path) {
446 final String target = basePath + "/_" + type + path;
447 if (ResourceUtil.getResourceNoException(target) != null) {
448 return target;
449 }
450 return basePath + path;
451 }
452
453 public void addMapping(final String index, final String docType, final String indexName) {
454 final FessConfig fessConfig = ComponentUtil.getFessConfig();
455
456 final GetMappingsResponse getMappingsResponse =
457 client.admin().indices().prepareGetMappings(indexName).execute().actionGet(fessConfig.getIndexIndicesTimeout());
458 final ImmutableOpenMap<String, MappingMetadata> indexMappings = getMappingsResponse.mappings().get(indexName);
459 if (indexMappings == null || !indexMappings.containsKey("properties")) {
460 String source = null;
461 final String mappingFile = getResourcePath(indexConfigPath, fessConfig.getFesenType(), "/" + index + "/" + docType + ".json");
462 try {
463 source = FileUtil.readUTF8(mappingFile);
464 } catch (final Exception e) {
465 logger.warn("{} is not found.", mappingFile, e);
466 }
467 try {
468 final AcknowledgedResponse putMappingResponse = client.admin().indices().preparePutMapping(indexName)
469 .setSource(source, XContentType.JSON).execute().actionGet(fessConfig.getIndexIndicesTimeout());
470 if (putMappingResponse.isAcknowledged()) {
471 logger.info("Created {}/{} mapping.", indexName, docType);
472 } else {
473 logger.warn("Failed to create {}/{} mapping.", indexName, docType);
474 }
475
476 final String dataPath = getResourcePath(indexConfigPath, fessConfig.getFesenType(), "/" + index + "/" + docType + ".bulk");
477 if (ResourceUtil.isExist(dataPath)) {
478 insertBulkData(fessConfig, indexName, dataPath);
479 }
480 split(fessConfig.getAppExtensionNames(), ",").of(stream -> stream.filter(StringUtil::isNotBlank).forEach(name -> {
481 final String bulkPath =
482 getResourcePath(indexConfigPath, fessConfig.getFesenType(), "/" + index + "/" + docType + "_" + name + ".bulk");
483 if (ResourceUtil.isExist(bulkPath)) {
484 insertBulkData(fessConfig, indexName, bulkPath);
485 }
486 }));
487 } catch (final Exception e) {
488 logger.warn("Failed to create {}/{} mapping.", indexName, docType, e);
489 }
490 } else if (logger.isDebugEnabled()) {
491 logger.debug("{}/{} mapping exists.", indexName, docType);
492 }
493 }
494
495 public boolean updateAlias(final String newIndex) {
496 final FessConfig fessConfig = ComponentUtil.getFessConfig();
497 final String updateAlias = fessConfig.getIndexDocumentUpdateIndex();
498 final String searchAlias = fessConfig.getIndexDocumentSearchIndex();
499 final GetIndexResponse response1 =
500 client.admin().indices().prepareGetIndex().addIndices(updateAlias).execute().actionGet(fessConfig.getIndexIndicesTimeout());
501 final String[] updateIndices = response1.indices();
502 final GetIndexResponse response2 =
503 client.admin().indices().prepareGetIndex().addIndices(searchAlias).execute().actionGet(fessConfig.getIndexIndicesTimeout());
504 final String[] searchIndices = response2.indices();
505
506 final IndicesAliasesRequestBuilder builder =
507 client.admin().indices().prepareAliases().addAlias(newIndex, updateAlias).addAlias(newIndex, searchAlias);
508 for (final String index : updateIndices) {
509 builder.removeAlias(index, updateAlias);
510 }
511 for (final String index : searchIndices) {
512 builder.removeAlias(index, searchAlias);
513 }
514 final AcknowledgedResponse response = builder.execute().actionGet(fessConfig.getIndexIndicesTimeout());
515 return response.isAcknowledged();
516 }
517
518 protected void createAlias(final String index, final String createdIndexName) {
519 final FessConfig fessConfig = ComponentUtil.getFessConfig();
520
521 final String aliasConfigDirPath = getResourcePath(indexConfigPath, fessConfig.getFesenType(), "/" + index + "/alias");
522 try {
523 final File aliasConfigDir = ResourceUtil.getResourceAsFile(aliasConfigDirPath);
524 if (aliasConfigDir.isDirectory()) {
525 stream(aliasConfigDir.listFiles((dir, name) -> name.endsWith(".json"))).of(stream -> stream.forEach(f -> {
526 final String aliasName = f.getName().replaceFirst(".json$", "");
527 String source = FileUtil.readUTF8(f);
528 if ("{}".equals(source.trim())) {
529 source = null;
530 }
531 final AcknowledgedResponse response = client.admin().indices().prepareAliases()
532 .addAlias(createdIndexName, aliasName, source).execute().actionGet(fessConfig.getIndexIndicesTimeout());
533 if (response.isAcknowledged()) {
534 logger.info("Created {} alias for {}", aliasName, createdIndexName);
535 } else if (logger.isDebugEnabled()) {
536 logger.debug("Failed to create {} alias for {}", aliasName, createdIndexName);
537 }
538 }));
539 }
540 } catch (final ResourceNotFoundRuntimeException e) {
541
542 } catch (final Exception e) {
543 logger.warn("{} is not found.", aliasConfigDirPath, e);
544 }
545 }
546
547 protected void sendConfigFiles(final String index) {
548 configListMap.getOrDefault(index, Collections.emptyList()).forEach(path -> {
549 String source = null;
550 final String filePath = indexConfigPath + "/" + index + "/" + path;
551 try {
552 source = FileUtil.readUTF8(filePath);
553 try (CurlResponse response =
554 ComponentUtil.getCurlHelper().post("/_configsync/file").param("path", path).body(source).execute()) {
555 if (response.getHttpStatusCode() == 200) {
556 logger.info("Register {} to {}", path, index);
557 } else if (response.getContentException() != null) {
558 logger.warn("Invalid request for {}.", path, response.getContentException());
559 } else {
560 logger.warn("Invalid request for {}. The response is {}", path, response.getContentAsString());
561 }
562 }
563 } catch (final Exception e) {
564 logger.warn("Failed to register {}", filePath, e);
565 }
566 });
567 try (CurlResponse response = ComponentUtil.getCurlHelper().post("/_configsync/flush").execute()) {
568 if (response.getHttpStatusCode() == 200) {
569 logger.info("Flushed config files.");
570 } else {
571 logger.warn("Failed to flush config files.");
572 }
573 } catch (final Exception e) {
574 logger.warn("Failed to flush config files.", e);
575 }
576 }
577
578 protected String generateNewIndexName(final String configIndex) {
579 return configIndex + "." + new SimpleDateFormat("yyyyMMdd").format(new Date());
580 }
581
582 protected void insertBulkData(final FessConfig fessConfig, final String configIndex, final String dataPath) {
583 try {
584 final BulkRequestBuilder builder = client.prepareBulk();
585 final ObjectMapper mapper = new ObjectMapper();
586 Arrays.stream(FileUtil.readUTF8(dataPath).split("\n")).reduce((prev, line) -> {
587 try {
588 if (StringUtil.isBlank(prev)) {
589 final Map<String, Map<String, String>> result =
590 mapper.readValue(line, new TypeReference<Map<String, Map<String, String>>>() {
591 });
592 if (result.containsKey("index") || result.containsKey("update")) {
593 return line;
594 }
595 if (result.containsKey("delete")) {
596 return StringUtil.EMPTY;
597 }
598 } else {
599 final Map<String, Map<String, String>> result =
600 mapper.readValue(prev, new TypeReference<Map<String, Map<String, String>>>() {
601 });
602 if (result.containsKey("index")) {
603 final IndexRequestBuilder requestBuilder = client.prepareIndex().setIndex(configIndex)
604 .setId(result.get("index").get("_id")).setSource(line, XContentType.JSON);
605 builder.add(requestBuilder);
606 }
607 }
608 } catch (final Exception e) {
609 logger.warn("Failed to parse {}", dataPath);
610 }
611 return StringUtil.EMPTY;
612 });
613 final BulkResponse response = builder.execute().actionGet(fessConfig.getIndexBulkTimeout());
614 if (response.hasFailures()) {
615 logger.warn("Failed to register {}: {}", dataPath, response.buildFailureMessage());
616 }
617 } catch (final Exception e) {
618 logger.warn("Failed to create {} mapping.", configIndex);
619 }
620 }
621
622 protected void waitForYellowStatus(final FessConfig fessConfig) {
623 Exception cause = null;
624 final long startTime = System.currentTimeMillis();
625 for (int i = 0; i < maxEsStatusRetry; i++) {
626 try {
627 final ClusterHealthResponse response = client.admin().cluster().prepareHealth().setWaitForYellowStatus().execute()
628 .actionGet(fessConfig.getIndexHealthTimeout());
629 if (logger.isDebugEnabled()) {
630 logger.debug("Fesen Cluster Status: {}", response.getStatus());
631 }
632 return;
633 } catch (final Exception e) {
634 cause = e;
635 }
636 if (cause instanceof FesenStatusException) {
637 final RestStatus status = ((FesenStatusException) cause).status();
638 switch (status) {
639 case UNAUTHORIZED:
640 logger.warn("[{}] Unauthorized access: {}", i, System.getProperty(Constants.FESS_ES_HTTP_ADDRESS), cause);
641 break;
642 default:
643 logger.debug("[{}][{}] Failed to access to Fesen ({})", i, status, System.getProperty(Constants.FESS_ES_HTTP_ADDRESS),
644 cause);
645 break;
646 }
647 } else if (logger.isDebugEnabled()) {
648 logger.debug("[{}] Failed to access to Fesen ({})", i, System.getProperty(Constants.FESS_ES_HTTP_ADDRESS), cause);
649 }
650 ThreadUtil.sleep(1000L);
651 }
652 final String message = "Fesen (" + System.getProperty(Constants.FESS_ES_HTTP_ADDRESS)
653 + ") is not available. Check the state of your Fesen cluster (" + clusterName + ") in "
654 + (System.currentTimeMillis() - startTime) + "ms.";
655 throw new ContainerInitFailureException(message, cause);
656 }
657
658 protected void waitForConfigSyncStatus() {
659 FessSystemException cause = null;
660 for (int i = 0; i < maxConfigSyncStatusRetry; i++) {
661 try (CurlResponse response = ComponentUtil.getCurlHelper().get("/_configsync/wait").param("status", "green").execute()) {
662 final int httpStatusCode = response.getHttpStatusCode();
663 if (httpStatusCode == 200) {
664 logger.info("ConfigSync is ready.");
665 return;
666 }
667 final String message = "Configsync is not available. HTTP Status is " + httpStatusCode;
668 if (response.getContentException() != null) {
669 throw new FessSystemException(message, response.getContentException());
670 }
671 throw new FessSystemException(message);
672 } catch (final Exception e) {
673 cause = new FessSystemException("Configsync is not available.", e);
674 }
675 if (logger.isDebugEnabled()) {
676 logger.debug("Failed to access to configsync:{}", i, cause);
677 }
678 ThreadUtil.sleep(1000L);
679 }
680 throw cause;
681 }
682
683 @Override
684 @PreDestroy
685 public void close() {
686 if (runner != null) {
687 try {
688 client.admin().indices().prepareFlush().setForce(true).execute()
689 .actionGet(ComponentUtil.getFessConfig().getIndexIndicesTimeout());
690 } catch (final Exception e) {
691 logger.warn("Failed to flush indices.", e);
692 }
693 }
694 try {
695 client.close();
696 } catch (final FesenException e) {
697 logger.warn("Failed to close Client: {}", client, e);
698 }
699 }
700
701 public long updateByQuery(final String index, final Function<SearchRequestBuilder, SearchRequestBuilder> option,
702 final BiFunction<UpdateRequestBuilder, SearchHit, UpdateRequestBuilder> builder) {
703
704 final FessConfig fessConfig = ComponentUtil.getFessConfig();
705 SearchResponse response = option.apply(client.prepareSearch(index).setScroll(scrollForUpdate).setSize(sizeForUpdate)
706 .setPreference(Constants.SEARCH_PREFERENCE_LOCAL)).execute().actionGet(fessConfig.getIndexScrollSearchTimeout());
707
708 int count = 0;
709 String scrollId = response.getScrollId();
710 try {
711 while (scrollId != null) {
712 final SearchHits searchHits = response.getHits();
713 final SearchHit[] hits = searchHits.getHits();
714 if (hits.length == 0) {
715 break;
716 }
717
718 final BulkRequestBuilder bulkRequest = client.prepareBulk();
719 for (final SearchHit hit : hits) {
720 final UpdateRequestBuilder requestBuilder =
721 builder.apply(client.prepareUpdate().setIndex(index).setId(hit.getId()), hit);
722 if (requestBuilder != null) {
723 bulkRequest.add(requestBuilder);
724 }
725 count++;
726 }
727 final BulkResponse bulkResponse = bulkRequest.execute().actionGet(fessConfig.getIndexBulkTimeout());
728 if (bulkResponse.hasFailures()) {
729 throw new IllegalBehaviorStateException(bulkResponse.buildFailureMessage());
730 }
731
732 response = client.prepareSearchScroll(scrollId).setScroll(scrollForUpdate).execute()
733 .actionGet(fessConfig.getIndexBulkTimeout());
734 if (!scrollId.equals(response.getScrollId())) {
735 deleteScrollContext(scrollId);
736 }
737 scrollId = response.getScrollId();
738 }
739 } finally {
740 deleteScrollContext(scrollId);
741 }
742 return count;
743 }
744
745 public long deleteByQuery(final String index, final QueryBuilder queryBuilder) {
746
747 final FessConfig fessConfig = ComponentUtil.getFessConfig();
748 SearchResponse response = client.prepareSearch(index).setScroll(scrollForDelete).setSize(sizeForDelete)
749 .setFetchSource(new String[] { fessConfig.getIndexFieldId() }, null).setQuery(queryBuilder)
750 .setPreference(Constants.SEARCH_PREFERENCE_LOCAL).execute().actionGet(fessConfig.getIndexScrollSearchTimeout());
751
752 int count = 0;
753 String scrollId = response.getScrollId();
754 try {
755 while (scrollId != null) {
756 final SearchHits searchHits = response.getHits();
757 final SearchHit[] hits = searchHits.getHits();
758 if (hits.length == 0) {
759 break;
760 }
761
762 final BulkRequestBuilder bulkRequest = client.prepareBulk();
763 for (final SearchHit hit : hits) {
764 bulkRequest.add(client.prepareDelete().setIndex(index).setId(hit.getId()));
765 count++;
766 }
767 final BulkResponse bulkResponse = bulkRequest.execute().actionGet(fessConfig.getIndexBulkTimeout());
768 if (bulkResponse.hasFailures()) {
769 throw new IllegalBehaviorStateException(bulkResponse.buildFailureMessage());
770 }
771
772 response = client.prepareSearchScroll(scrollId).setScroll(scrollForDelete).execute()
773 .actionGet(fessConfig.getIndexBulkTimeout());
774 if (!scrollId.equals(response.getScrollId())) {
775 deleteScrollContext(scrollId);
776 }
777 scrollId = response.getScrollId();
778 }
779 } finally {
780 deleteScrollContext(scrollId);
781 }
782 return count;
783 }
784
785 protected void deleteScrollContext(final String scrollId) {
786 if (scrollId != null) {
787 client.prepareClearScroll().addScrollId(scrollId)
788 .execute(ActionListener.wrap(res -> {}, e -> logger.warn("Failed to clear the scroll context.", e)));
789 }
790 }
791
792 protected <T> T get(final String index, final String type, final String id, final SearchCondition<GetRequestBuilder> condition,
793 final SearchResult<T, GetRequestBuilder, GetResponse> searchResult) {
794 final long startTime = System.currentTimeMillis();
795
796 GetResponse response = null;
797 final GetRequestBuilder requestBuilder = client.prepareGet(index, type, id);
798 if (condition.build(requestBuilder)) {
799 response = requestBuilder.execute().actionGet(ComponentUtil.getFessConfig().getIndexSearchTimeout());
800 }
801 final long execTime = System.currentTimeMillis() - startTime;
802
803 return searchResult.build(requestBuilder, execTime, OptionalEntity.ofNullable(response, () -> {}));
804 }
805
806 public <T> T search(final String index, final SearchCondition<SearchRequestBuilder> condition,
807 final SearchResult<T, SearchRequestBuilder, SearchResponse> searchResult) {
808 final long startTime = System.currentTimeMillis();
809
810 SearchResponse searchResponse = null;
811 final SearchRequestBuilder searchRequestBuilder = client.prepareSearch(index);
812 if (condition.build(searchRequestBuilder)) {
813
814 final FessConfig fessConfig = ComponentUtil.getFessConfig();
815 final long queryTimeout = fessConfig.getQueryTimeoutAsInteger().longValue();
816 if (queryTimeout >= 0) {
817 searchRequestBuilder.setTimeout(TimeValue.timeValueMillis(queryTimeout));
818 }
819
820 try {
821 if (logger.isDebugEnabled()) {
822 logger.debug("Query DSL:\n{}", searchRequestBuilder);
823 }
824 searchResponse = searchRequestBuilder.execute().actionGet(ComponentUtil.getFessConfig().getIndexSearchTimeout());
825 } catch (final SearchPhaseExecutionException e) {
826 throw new InvalidQueryException(messages -> messages.addErrorsInvalidQueryParseError(UserMessages.GLOBAL_PROPERTY_KEY),
827 "Invalid query: " + searchRequestBuilder, e);
828 }
829 }
830 final long execTime = System.currentTimeMillis() - startTime;
831
832 return searchResult.build(searchRequestBuilder, execTime, OptionalEntity.ofNullable(searchResponse, () -> {}));
833 }
834
835 public <T> long scrollSearch(final String index, final SearchCondition<SearchRequestBuilder> condition,
836 final EntityCreator<T, SearchResponse, SearchHit> creator, final BooleanFunction<T> cursor) {
837 long count = 0;
838
839 final SearchRequestBuilder searchRequestBuilder = client.prepareSearch(index).setScroll(scrollForSearch);
840 if (condition.build(searchRequestBuilder)) {
841 final FessConfig fessConfig = ComponentUtil.getFessConfig();
842
843 String scrollId = null;
844 try {
845 if (logger.isDebugEnabled()) {
846 logger.debug("Query DSL:\n{}", searchRequestBuilder);
847 }
848 SearchResponse response = searchRequestBuilder.execute().actionGet(ComponentUtil.getFessConfig().getIndexSearchTimeout());
849
850 scrollId = response.getScrollId();
851 while (scrollId != null) {
852 final SearchHits searchHits = response.getHits();
853 final SearchHit[] hits = searchHits.getHits();
854 if (hits.length == 0) {
855 break;
856 }
857
858 for (final SearchHit hit : hits) {
859 count++;
860 if (!cursor.apply(creator.build(response, hit))) {
861 break;
862 }
863 }
864
865 response = client.prepareSearchScroll(scrollId).setScroll(scrollForDelete).execute()
866 .actionGet(fessConfig.getIndexBulkTimeout());
867 if (!scrollId.equals(response.getScrollId())) {
868 deleteScrollContext(scrollId);
869 }
870 scrollId = response.getScrollId();
871 }
872 } catch (final SearchPhaseExecutionException e) {
873 throw new InvalidQueryException(messages -> messages.addErrorsInvalidQueryParseError(UserMessages.GLOBAL_PROPERTY_KEY),
874 "Invalid query: " + searchRequestBuilder, e);
875 } finally {
876 deleteScrollContext(scrollId);
877 }
878 }
879
880 return count;
881 }
882
883 public OptionalEntity<Map<String, Object>> getDocument(final String index, final SearchCondition<SearchRequestBuilder> condition) {
884 return getDocument(index, condition, (response, hit) -> {
885 final FessConfig fessConfig = ComponentUtil.getFessConfig();
886 final Map<String, Object> source = hit.getSourceAsMap();
887 if (source != null) {
888 final Map<String, Object> docMap = new HashMap<>(source);
889 docMap.put(fessConfig.getIndexFieldId(), hit.getId());
890 docMap.put(fessConfig.getIndexFieldVersion(), hit.getVersion());
891 docMap.put(fessConfig.getIndexFieldSeqNo(), hit.getSeqNo());
892 docMap.put(fessConfig.getIndexFieldPrimaryTerm(), hit.getPrimaryTerm());
893 return docMap;
894 }
895 final Map<String, DocumentField> fields = hit.getFields();
896 if (fields != null) {
897 final Map<String, Object> docMap = fields.entrySet().stream()
898 .collect(Collectors.toMap(Entry<String, DocumentField>::getKey, e -> (Object) e.getValue().getValues()));
899 docMap.put(fessConfig.getIndexFieldId(), hit.getId());
900 docMap.put(fessConfig.getIndexFieldVersion(), hit.getVersion());
901 docMap.put(fessConfig.getIndexFieldSeqNo(), hit.getSeqNo());
902 docMap.put(fessConfig.getIndexFieldPrimaryTerm(), hit.getPrimaryTerm());
903 return docMap;
904 }
905 return null;
906 });
907 }
908
909 protected <T> OptionalEntity<T> getDocument(final String index, final SearchCondition<SearchRequestBuilder> condition,
910 final EntityCreator<T, SearchResponse, SearchHit> creator) {
911 return search(index, searchRequestBuilder -> {
912 searchRequestBuilder.setVersion(true);
913 return condition.build(searchRequestBuilder);
914 }, (queryBuilder, execTime, searchResponse) -> searchResponse.map(response -> {
915 final SearchHit[] hits = response.getHits().getHits();
916 if (hits.length > 0) {
917 return creator.build(response, hits[0]);
918 }
919 return null;
920 }));
921 }
922
923 public List<Map<String, Object>> getDocumentList(final String index, final SearchCondition<SearchRequestBuilder> condition) {
924 return getDocumentList(index, condition, (response, hit) -> {
925 final FessConfig fessConfig = ComponentUtil.getFessConfig();
926 final Map<String, Object> source = hit.getSourceAsMap();
927 if (source != null) {
928 final Map<String, Object> docMap = new HashMap<>(source);
929 docMap.put(fessConfig.getIndexFieldId(), hit.getId());
930 return docMap;
931 }
932 final Map<String, DocumentField> fields = hit.getFields();
933 if (fields != null) {
934 final Map<String, Object> docMap = fields.entrySet().stream()
935 .collect(Collectors.toMap(Entry<String, DocumentField>::getKey, e -> (Object) e.getValue().getValues()));
936 docMap.put(fessConfig.getIndexFieldId(), hit.getId());
937 return docMap;
938 }
939 return null;
940 });
941 }
942
943 protected <T> List<T> getDocumentList(final String index, final SearchCondition<SearchRequestBuilder> condition,
944 final EntityCreator<T, SearchResponse, SearchHit> creator) {
945 return search(index, condition, (searchRequestBuilder, execTime, searchResponse) -> {
946 final List<T> list = new ArrayList<>();
947 searchResponse.ifPresent(response -> response.getHits().forEach(hit -> {
948 list.add(creator.build(response, hit));
949 }));
950 return list;
951 });
952 }
953
954 public boolean update(final String index, final String id, final String field, final Object value) {
955 try {
956 final Result result = client.prepareUpdate().setIndex(index).setId(id).setDoc(field, value).execute()
957 .actionGet(ComponentUtil.getFessConfig().getIndexIndexTimeout()).getResult();
958 return result == Result.CREATED || result == Result.UPDATED;
959 } catch (final FesenException e) {
960 throw new SearchEngineClientException("Failed to set " + value + " to " + field + " for doc " + id, e);
961 }
962 }
963
964 public void refresh(final String... indices) {
965 client.admin().indices().prepareRefresh(indices).execute(new ActionListener<RefreshResponse>() {
966 @Override
967 public void onResponse(final RefreshResponse response) {
968 if (logger.isDebugEnabled()) {
969 logger.debug(() -> "Refreshed " + stream(indices).get(stream -> stream.collect(Collectors.joining(", "))));
970 }
971 }
972
973 @Override
974 public void onFailure(final Exception e) {
975 logger.error(() -> "Failed to refresh " + stream(indices).get(stream -> stream.collect(Collectors.joining(", "))), e);
976 }
977 });
978
979 }
980
981 public void flush(final String... indices) {
982 client.admin().indices().prepareFlush(indices).execute(new ActionListener<FlushResponse>() {
983
984 @Override
985 public void onResponse(final FlushResponse response) {
986 if (logger.isDebugEnabled()) {
987 logger.debug(() -> "Flushed " + stream(indices).get(stream -> stream.collect(Collectors.joining(", "))));
988 }
989 }
990
991 @Override
992 public void onFailure(final Exception e) {
993 logger.error(() -> "Failed to flush " + stream(indices).get(stream -> stream.collect(Collectors.joining(", "))), e);
994 }
995 });
996
997 }
998
999 public PingResponse ping() {
1000 try {
1001 final ClusterHealthResponse response =
1002 client.admin().cluster().prepareHealth().execute().actionGet(ComponentUtil.getFessConfig().getIndexHealthTimeout());
1003 return new PingResponse(response);
1004 } catch (final FesenException e) {
1005 throw new SearchEngineClientException("Failed to process a ping request.", e);
1006 }
1007 }
1008
1009 public void addAll(final String index, final List<Map<String, Object>> docList,
1010 final BiConsumer<Map<String, Object>, IndexRequestBuilder> options) {
1011 final FessConfig fessConfig = ComponentUtil.getFessConfig();
1012 final BulkRequestBuilder bulkRequestBuilder = client.prepareBulk();
1013 for (final Map<String, Object> doc : docList) {
1014 final Object id = doc.remove(fessConfig.getIndexFieldId());
1015 final IndexRequestBuilder builder = client.prepareIndex().setIndex(index).setId(id.toString()).setSource(new DocMap(doc));
1016 options.accept(doc, builder);
1017 bulkRequestBuilder.add(builder);
1018 }
1019 final BulkResponse response = bulkRequestBuilder.execute().actionGet(ComponentUtil.getFessConfig().getIndexBulkTimeout());
1020 if (response.hasFailures()) {
1021 if (logger.isDebugEnabled()) {
1022 final List<DocWriteRequest<?>> requests = bulkRequestBuilder.request().requests();
1023 final BulkItemResponse[] items = response.getItems();
1024 if (requests.size() == items.length) {
1025 for (int i = 0; i < requests.size(); i++) {
1026 final BulkItemResponse resp = items[i];
1027 if (resp.isFailed() && resp.getFailure() != null) {
1028 final DocWriteRequest<?> req = requests.get(i);
1029 final Failure failure = resp.getFailure();
1030 logger.debug("Failed Request: {}\n=>{}", req, failure.getMessage());
1031 }
1032 }
1033 }
1034 }
1035 throw new SearchEngineClientException(response.buildFailureMessage());
1036 }
1037 }
1038
1039 public static class SearchConditionBuilder {
1040 protected final SearchRequestBuilder searchRequestBuilder;
1041 protected String query;
1042 protected String[] responseFields;
1043 protected int offset = Constants.DEFAULT_START_COUNT;
1044 protected int size = Constants.DEFAULT_PAGE_SIZE;
1045 protected GeoInfo geoInfo;
1046 protected FacetInfo facetInfo;
1047 protected HighlightInfo highlightInfo;
1048 protected String similarDocHash;
1049 protected SearchRequestType searchRequestType = SearchRequestType.SEARCH;
1050 protected boolean isScroll = false;
1051 protected String trackTotalHits = null;
1052
1053 public static SearchConditionBuilder builder(final SearchRequestBuilder searchRequestBuilder) {
1054 return new SearchConditionBuilder(searchRequestBuilder);
1055 }
1056
1057 SearchConditionBuilder(final SearchRequestBuilder searchRequestBuilder) {
1058 this.searchRequestBuilder = searchRequestBuilder;
1059 }
1060
1061 public Map<String, Object> condition() {
1062 final Map<String, Object> params = new HashMap<>();
1063 params.put("query", query);
1064 params.put("responseFields", responseFields);
1065 params.put("offset", offset);
1066 params.put("size", size);
1067
1068
1069
1070 params.put("similarDocHash", similarDocHash);
1071 return params;
1072 }
1073
1074 public SearchConditionBuilder query(final String query) {
1075 this.query = query;
1076 return this;
1077 }
1078
1079 public SearchConditionBuilder searchRequestType(final SearchRequestType searchRequestType) {
1080 this.searchRequestType = searchRequestType;
1081 return this;
1082 }
1083
1084 public SearchConditionBuilder responseFields(final String[] responseFields) {
1085 this.responseFields = responseFields;
1086 return this;
1087 }
1088
1089 public SearchConditionBuilder offset(final int offset) {
1090 this.offset = offset;
1091 return this;
1092 }
1093
1094 public SearchConditionBuilder size(final int size) {
1095 this.size = size;
1096 return this;
1097 }
1098
1099 public SearchConditionBuilder geoInfo(final GeoInfo geoInfo) {
1100 this.geoInfo = geoInfo;
1101 return this;
1102 }
1103
1104 public SearchConditionBuilder highlightInfo(final HighlightInfo highlightInfo) {
1105 this.highlightInfo = highlightInfo;
1106 return this;
1107 }
1108
1109 public SearchConditionBuilder similarDocHash(final String similarDocHash) {
1110 if (StringUtil.isNotBlank(similarDocHash)) {
1111 this.similarDocHash = similarDocHash;
1112 }
1113 return this;
1114 }
1115
1116 public SearchConditionBuilder facetInfo(final FacetInfo facetInfo) {
1117 this.facetInfo = facetInfo;
1118 return this;
1119 }
1120
1121 public SearchConditionBuilder scroll() {
1122 this.isScroll = true;
1123 return this;
1124 }
1125
1126 public SearchConditionBuilder trackTotalHits(final String trackTotalHits) {
1127 this.trackTotalHits = trackTotalHits;
1128 return this;
1129 }
1130
1131 public boolean build() {
1132 if (StringUtil.isBlank(query)) {
1133 return false;
1134 }
1135
1136 final QueryHelper queryHelper = ComponentUtil.getQueryHelper();
1137 final FessConfig fessConfig = ComponentUtil.getFessConfig();
1138
1139 if (offset > fessConfig.getQueryMaxSearchResultOffsetAsInteger()) {
1140 throw new ResultOffsetExceededException("The number of result size is exceeded.");
1141 }
1142
1143 final QueryContext queryContext = buildQueryContext(queryHelper, fessConfig);
1144
1145 searchRequestBuilder.setFrom(offset).setSize(size);
1146
1147 buildTrackTotalHits(fessConfig);
1148
1149 if (responseFields != null) {
1150 searchRequestBuilder.setFetchSource(responseFields, null);
1151 }
1152
1153
1154 buildRescorer(queryHelper, fessConfig);
1155
1156
1157 buildSort(queryContext, fessConfig);
1158
1159
1160 if (highlightInfo != null) {
1161 buildHighlighter(queryHelper, fessConfig);
1162 }
1163
1164
1165 if (facetInfo != null) {
1166 buildFacet(queryHelper, fessConfig);
1167 }
1168
1169 if (!SearchRequestType.ADMIN_SEARCH.equals(searchRequestType) && !isScroll && fessConfig.isResultCollapsed()
1170 && similarDocHash == null) {
1171 searchRequestBuilder.setCollapse(getCollapseBuilder(fessConfig));
1172 }
1173
1174 searchRequestBuilder.setQuery(queryContext.getQueryBuilder());
1175 return true;
1176 }
1177
1178 protected void buildTrackTotalHits(final FessConfig fessConfig) {
1179 if (StringUtil.isNotBlank(trackTotalHits)) {
1180 if (Constants.TRUE.equalsIgnoreCase(trackTotalHits) || Constants.FALSE.equalsIgnoreCase(trackTotalHits)) {
1181 searchRequestBuilder.setTrackTotalHits(Boolean.parseBoolean(trackTotalHits));
1182 return;
1183 }
1184 try {
1185 searchRequestBuilder.setTrackTotalHitsUpTo(Integer.parseInt(trackTotalHits));
1186 return;
1187 } catch (final NumberFormatException e) {
1188
1189 }
1190 }
1191 final Object trackTotalHitsValue = fessConfig.getQueryTrackTotalHitsValue();
1192 if (trackTotalHitsValue instanceof Boolean) {
1193 searchRequestBuilder.setTrackTotalHits((Boolean) trackTotalHitsValue);
1194 } else if (trackTotalHitsValue instanceof Number) {
1195 searchRequestBuilder.setTrackTotalHitsUpTo(((Number) trackTotalHitsValue).intValue());
1196 }
1197 }
1198
1199 protected void buildFacet(final QueryHelper queryHelper, final FessConfig fessConfig) {
1200 stream(facetInfo.field).of(stream -> stream.forEach(f -> {
1201 if (!queryHelper.isFacetField(f)) {
1202 throw new SearchQueryException("Invalid facet field: " + f);
1203 }
1204 final String encodedField = BaseEncoding.base64().encode(f.getBytes(StandardCharsets.UTF_8));
1205 final TermsAggregationBuilder termsBuilder =
1206 AggregationBuilders.terms(Constants.FACET_FIELD_PREFIX + encodedField).field(f);
1207 termsBuilder.order(facetInfo.getBucketOrder());
1208 if (facetInfo.size != null) {
1209 termsBuilder.size(facetInfo.size);
1210 }
1211 if (facetInfo.minDocCount != null) {
1212 termsBuilder.minDocCount(facetInfo.minDocCount);
1213 }
1214 if (facetInfo.missing != null) {
1215 termsBuilder.missing(facetInfo.missing);
1216 }
1217 searchRequestBuilder.addAggregation(termsBuilder);
1218 }));
1219 stream(facetInfo.query).of(stream -> stream.forEach(fq -> {
1220 final QueryContext facetContext = new QueryContext(fq, false);
1221 queryHelper.buildBaseQuery(facetContext, c -> {});
1222 final String encodedFacetQuery = BaseEncoding.base64().encode(fq.getBytes(StandardCharsets.UTF_8));
1223 final FilterAggregationBuilder filterBuilder =
1224 AggregationBuilders.filter(Constants.FACET_QUERY_PREFIX + encodedFacetQuery, facetContext.getQueryBuilder());
1225 searchRequestBuilder.addAggregation(filterBuilder);
1226 }));
1227 }
1228
1229 protected void buildHighlighter(final QueryHelper queryHelper, final FessConfig fessConfig) {
1230 final String highlighterType = highlightInfo.getType();
1231 final int fragmentSize = highlightInfo.getFragmentSize();
1232 final int numOfFragments = highlightInfo.getNumOfFragments();
1233 final int fragmentOffset = highlightInfo.getFragmentOffset();
1234 final char[] boundaryChars = fessConfig.getQueryHighlightBoundaryCharsAsArray();
1235 final int boundaryMaxScan = fessConfig.getQueryHighlightBoundaryMaxScanAsInteger();
1236 final String boundaryScannerType = fessConfig.getQueryHighlightBoundaryScanner();
1237 final boolean forceSource = fessConfig.isQueryHighlightForceSource();
1238 final String fragmenter = fessConfig.getQueryHighlightFragmenter();
1239 final int noMatchSize = fessConfig.getQueryHighlightNoMatchSizeAsInteger();
1240 final String order = fessConfig.getQueryHighlightOrder();
1241 final int phraseLimit = fessConfig.getQueryHighlightPhraseLimitAsInteger();
1242 final String encoder = fessConfig.getQueryHighlightEncoder();
1243 final HighlightBuilder highlightBuilder = new HighlightBuilder();
1244 queryHelper.highlightedFields(stream -> stream.forEach(hf -> highlightBuilder
1245 .field(new HighlightBuilder.Field(hf).highlighterType(highlighterType).fragmentSize(fragmentSize)
1246 .numOfFragments(numOfFragments).boundaryChars(boundaryChars).boundaryMaxScan(boundaryMaxScan)
1247 .boundaryScannerType(boundaryScannerType).forceSource(forceSource).fragmenter(fragmenter)
1248 .fragmentOffset(fragmentOffset).noMatchSize(noMatchSize).order(order).phraseLimit(phraseLimit))
1249 .encoder(encoder)));
1250 searchRequestBuilder.highlighter(highlightBuilder);
1251 }
1252
1253 protected void buildSort(final QueryContext queryContext, final FessConfig fessConfig) {
1254 queryContext.sortBuilders().forEach(sortBuilder -> searchRequestBuilder.addSort(sortBuilder));
1255 }
1256
1257 protected void buildRescorer(final QueryHelper queryHelper, final FessConfig fessConfig) {
1258 stream(queryHelper.getRescorers(condition())).of(stream -> stream.forEach(searchRequestBuilder::addRescorer));
1259 }
1260
1261 protected QueryContext buildQueryContext(final QueryHelper queryHelper, final FessConfig fessConfig) {
1262 return queryHelper.build(searchRequestType, query, context -> {
1263 if (SearchRequestType.ADMIN_SEARCH.equals(searchRequestType)) {
1264 context.skipRoleQuery();
1265 } else if (similarDocHash != null) {
1266 final DocumentHelper documentHelper = ComponentUtil.getDocumentHelper();
1267 context.addQuery(boolQuery -> {
1268 boolQuery.filter(QueryBuilders.termQuery(fessConfig.getIndexFieldContentMinhashBits(),
1269 documentHelper.decodeSimilarDocHash(similarDocHash)));
1270 });
1271 }
1272
1273 if (geoInfo != null && geoInfo.toQueryBuilder() != null) {
1274 context.addQuery(boolQuery -> boolQuery.filter(geoInfo.toQueryBuilder()));
1275 }
1276 });
1277 }
1278
1279 protected CollapseBuilder getCollapseBuilder(final FessConfig fessConfig) {
1280 final InnerHitBuilder innerHitBuilder = new InnerHitBuilder().setName(fessConfig.getQueryCollapseInnerHitsName())
1281 .setSize(fessConfig.getQueryCollapseInnerHitsSizeAsInteger());
1282 fessConfig.getQueryCollapseInnerHitsSortBuilders()
1283 .ifPresent(builders -> stream(builders).of(stream -> stream.forEach(innerHitBuilder::addSort)));
1284 return new CollapseBuilder(fessConfig.getIndexFieldContentMinhashBits())
1285 .setMaxConcurrentGroupRequests(fessConfig.getQueryCollapseMaxConcurrentGroupResultsAsInteger())
1286 .setInnerHits(innerHitBuilder);
1287 }
1288 }
1289
1290 public boolean store(final String index, final Object obj) {
1291 final FessConfig fessConfig = ComponentUtil.getFessConfig();
1292 @SuppressWarnings("unchecked")
1293 final Map<String, Object> source = obj instanceof Map ? (Map<String, Object>) obj : BeanUtil.copyBeanToNewMap(obj);
1294 final String id = (String) source.remove(fessConfig.getIndexFieldId());
1295 source.remove(fessConfig.getIndexFieldVersion());
1296 final Number seqNo = (Number) source.remove(fessConfig.getIndexFieldSeqNo());
1297 final Number primaryTerm = (Number) source.remove(fessConfig.getIndexFieldPrimaryTerm());
1298 IndexResponse response;
1299 try {
1300 if (id == null) {
1301
1302
1303 response = client.prepareIndex().setIndex(index).setSource(new DocMap(source)).setRefreshPolicy(RefreshPolicy.IMMEDIATE)
1304 .setOpType(OpType.CREATE).execute().actionGet(fessConfig.getIndexIndexTimeout());
1305 } else {
1306
1307 final IndexRequestBuilder builder = client.prepareIndex().setIndex(index).setId(id).setSource(new DocMap(source))
1308 .setRefreshPolicy(RefreshPolicy.IMMEDIATE).setOpType(OpType.INDEX);
1309 if (seqNo != null) {
1310 builder.setIfSeqNo(seqNo.longValue());
1311 }
1312 if (primaryTerm != null) {
1313 builder.setIfPrimaryTerm(primaryTerm.longValue());
1314 }
1315 response = builder.execute().actionGet(fessConfig.getIndexIndexTimeout());
1316 }
1317 final Result result = response.getResult();
1318 return result == Result.CREATED || result == Result.UPDATED;
1319 } catch (final FesenException e) {
1320 throw new SearchEngineClientException("Failed to store: " + obj, e);
1321 }
1322 }
1323
1324 public boolean delete(final String index, final String id) {
1325 return delete(index, id, null, null);
1326 }
1327
1328 public boolean delete(final String index, final String id, final Number seqNo, final Number primaryTerm) {
1329 try {
1330 final DeleteRequestBuilder builder = client.prepareDelete().setIndex(index).setId(id).setRefreshPolicy(RefreshPolicy.IMMEDIATE);
1331 if (seqNo != null) {
1332 builder.setIfSeqNo(seqNo.longValue());
1333 }
1334 if (primaryTerm != null) {
1335 builder.setIfPrimaryTerm(primaryTerm.longValue());
1336 }
1337 final DeleteResponse response = builder.execute().actionGet(ComponentUtil.getFessConfig().getIndexDeleteTimeout());
1338 return response.getResult() == Result.DELETED;
1339 } catch (final FesenException e) {
1340 throw new SearchEngineClientException("Failed to delete: " + index + "/" + id + "@" + seqNo + ":" + primaryTerm, e);
1341 }
1342 }
1343
1344 public void setIndexConfigPath(final String indexConfigPath) {
1345 this.indexConfigPath = indexConfigPath;
1346 }
1347
1348 public interface SearchCondition<B> {
1349 boolean build(B requestBuilder);
1350 }
1351
1352 public interface SearchResult<T, B, R> {
1353 T build(B requestBuilder, long execTime, OptionalEntity<R> response);
1354 }
1355
1356 public interface EntityCreator<T, R, H> {
1357 T build(R response, H hit);
1358 }
1359
1360 public void setClusterName(final String clusterName) {
1361 this.clusterName = clusterName;
1362 }
1363
1364
1365
1366
1367
1368 @Override
1369 public ThreadPool threadPool() {
1370 return client.threadPool();
1371 }
1372
1373 @Override
1374 public AdminClient admin() {
1375 return client.admin();
1376 }
1377
1378 @Override
1379 public ActionFuture<IndexResponse> index(final IndexRequest request) {
1380 return client.index(request);
1381 }
1382
1383 @Override
1384 public void index(final IndexRequest request, final ActionListener<IndexResponse> listener) {
1385 client.index(request, listener);
1386 }
1387
1388 @Override
1389 public IndexRequestBuilder prepareIndex() {
1390 return client.prepareIndex();
1391 }
1392
1393 @Override
1394 public ActionFuture<UpdateResponse> update(final UpdateRequest request) {
1395 return client.update(request);
1396 }
1397
1398 @Override
1399 public void update(final UpdateRequest request, final ActionListener<UpdateResponse> listener) {
1400 client.update(request, listener);
1401 }
1402
1403 @Override
1404 public UpdateRequestBuilder prepareUpdate() {
1405 return client.prepareUpdate();
1406 }
1407
1408 @Override
1409 public UpdateRequestBuilder prepareUpdate(final String index, final String type, final String id) {
1410 return client.prepareUpdate(index, type, id);
1411 }
1412
1413 @Override
1414 public IndexRequestBuilder prepareIndex(final String index, final String type) {
1415 return client.prepareIndex(index, type);
1416 }
1417
1418 @Override
1419 public IndexRequestBuilder prepareIndex(final String index, final String type, final String id) {
1420 return client.prepareIndex(index, type, id);
1421 }
1422
1423 @Override
1424 public ActionFuture<DeleteResponse> delete(final DeleteRequest request) {
1425 return client.delete(request);
1426 }
1427
1428 @Override
1429 public void delete(final DeleteRequest request, final ActionListener<DeleteResponse> listener) {
1430 client.delete(request, listener);
1431 }
1432
1433 @Override
1434 public DeleteRequestBuilder prepareDelete() {
1435 return client.prepareDelete();
1436 }
1437
1438 @Override
1439 public DeleteRequestBuilder prepareDelete(final String index, final String type, final String id) {
1440 return client.prepareDelete(index, type, id);
1441 }
1442
1443 @Override
1444 public ActionFuture<BulkResponse> bulk(final BulkRequest request) {
1445 return client.bulk(request);
1446 }
1447
1448 @Override
1449 public void bulk(final BulkRequest request, final ActionListener<BulkResponse> listener) {
1450 client.bulk(request, listener);
1451 }
1452
1453 @Override
1454 public BulkRequestBuilder prepareBulk() {
1455 return client.prepareBulk();
1456 }
1457
1458 @Override
1459 public ActionFuture<GetResponse> get(final GetRequest request) {
1460 return client.get(request);
1461 }
1462
1463 @Override
1464 public void get(final GetRequest request, final ActionListener<GetResponse> listener) {
1465 client.get(request, listener);
1466 }
1467
1468 @Override
1469 public GetRequestBuilder prepareGet() {
1470 return client.prepareGet();
1471 }
1472
1473 @Override
1474 public GetRequestBuilder prepareGet(final String index, final String type, final String id) {
1475 return client.prepareGet(index, type, id);
1476 }
1477
1478 @Override
1479 public ActionFuture<MultiGetResponse> multiGet(final MultiGetRequest request) {
1480 return client.multiGet(request);
1481 }
1482
1483 @Override
1484 public void multiGet(final MultiGetRequest request, final ActionListener<MultiGetResponse> listener) {
1485 client.multiGet(request, listener);
1486 }
1487
1488 @Override
1489 public MultiGetRequestBuilder prepareMultiGet() {
1490 return client.prepareMultiGet();
1491 }
1492
1493 @Override
1494 public ActionFuture<SearchResponse> search(final SearchRequest request) {
1495 return client.search(request);
1496 }
1497
1498 @Override
1499 public void search(final SearchRequest request, final ActionListener<SearchResponse> listener) {
1500 client.search(request, listener);
1501 }
1502
1503 @Override
1504 public SearchRequestBuilder prepareSearch(final String... indices) {
1505 return client.prepareSearch(indices);
1506 }
1507
1508 @Override
1509 public ActionFuture<SearchResponse> searchScroll(final SearchScrollRequest request) {
1510 return client.searchScroll(request);
1511 }
1512
1513 @Override
1514 public void searchScroll(final SearchScrollRequest request, final ActionListener<SearchResponse> listener) {
1515 client.searchScroll(request, listener);
1516 }
1517
1518 @Override
1519 public SearchScrollRequestBuilder prepareSearchScroll(final String scrollId) {
1520 return client.prepareSearchScroll(scrollId);
1521 }
1522
1523 @Override
1524 public ActionFuture<MultiSearchResponse> multiSearch(final MultiSearchRequest request) {
1525 return client.multiSearch(request);
1526 }
1527
1528 @Override
1529 public void multiSearch(final MultiSearchRequest request, final ActionListener<MultiSearchResponse> listener) {
1530 client.multiSearch(request, listener);
1531 }
1532
1533 @Override
1534 public MultiSearchRequestBuilder prepareMultiSearch() {
1535 return client.prepareMultiSearch();
1536 }
1537
1538 @Override
1539 public ExplainRequestBuilder prepareExplain(final String index, final String type, final String id) {
1540 return client.prepareExplain(index, type, id);
1541 }
1542
1543 @Override
1544 public ActionFuture<ExplainResponse> explain(final ExplainRequest request) {
1545 return client.explain(request);
1546 }
1547
1548 @Override
1549 public void explain(final ExplainRequest request, final ActionListener<ExplainResponse> listener) {
1550 client.explain(request, listener);
1551 }
1552
1553 @Override
1554 public ClearScrollRequestBuilder prepareClearScroll() {
1555 return client.prepareClearScroll();
1556 }
1557
1558 @Override
1559 public ActionFuture<ClearScrollResponse> clearScroll(final ClearScrollRequest request) {
1560 return client.clearScroll(request);
1561 }
1562
1563 @Override
1564 public void clearScroll(final ClearScrollRequest request, final ActionListener<ClearScrollResponse> listener) {
1565 client.clearScroll(request, listener);
1566 }
1567
1568 @Override
1569 public Settings settings() {
1570 return client.settings();
1571 }
1572
1573 @Override
1574 public ActionFuture<TermVectorsResponse> termVectors(final TermVectorsRequest request) {
1575 return client.termVectors(request);
1576 }
1577
1578 @Override
1579 public void termVectors(final TermVectorsRequest request, final ActionListener<TermVectorsResponse> listener) {
1580 client.termVectors(request, listener);
1581 }
1582
1583 @Override
1584 public TermVectorsRequestBuilder prepareTermVectors() {
1585 return client.prepareTermVectors();
1586 }
1587
1588 @Override
1589 public TermVectorsRequestBuilder prepareTermVectors(final String index, final String type, final String id) {
1590 return client.prepareTermVectors(index, type, id);
1591 }
1592
1593 @Override
1594 public ActionFuture<MultiTermVectorsResponse> multiTermVectors(final MultiTermVectorsRequest request) {
1595 return client.multiTermVectors(request);
1596 }
1597
1598 @Override
1599 public void multiTermVectors(final MultiTermVectorsRequest request, final ActionListener<MultiTermVectorsResponse> listener) {
1600 client.multiTermVectors(request, listener);
1601 }
1602
1603 @Override
1604 public MultiTermVectorsRequestBuilder prepareMultiTermVectors() {
1605 return client.prepareMultiTermVectors();
1606 }
1607
1608 public void setSizeForUpdate(final int sizeForUpdate) {
1609 this.sizeForUpdate = sizeForUpdate;
1610 }
1611
1612 public void setScrollForUpdate(final String scrollForUpdate) {
1613 this.scrollForUpdate = scrollForUpdate;
1614 }
1615
1616 public void setSizeForDelete(final int sizeForDelete) {
1617 this.sizeForDelete = sizeForDelete;
1618 }
1619
1620 public void setScrollForDelete(final String scrollForDelete) {
1621 this.scrollForDelete = scrollForDelete;
1622 }
1623
1624 public void setScrollForSearch(final String scrollForSearch) {
1625 this.scrollForSearch = scrollForSearch;
1626 }
1627
1628 public void setMaxConfigSyncStatusRetry(final int maxConfigSyncStatusRetry) {
1629 this.maxConfigSyncStatusRetry = maxConfigSyncStatusRetry;
1630 }
1631
1632 public void setMaxEsStatusRetry(final int maxEsStatusRetry) {
1633 this.maxEsStatusRetry = maxEsStatusRetry;
1634 }
1635
1636 @Override
1637 public Client filterWithHeader(final Map<String, String> headers) {
1638 return client.filterWithHeader(headers);
1639 }
1640
1641 @Override
1642 public <Request extends ActionRequest, Response extends ActionResponse> ActionFuture<Response> execute(
1643 final ActionType<Response> action, final Request request) {
1644 return client.execute(action, request);
1645 }
1646
1647 @Override
1648 public <Request extends ActionRequest, Response extends ActionResponse> void execute(final ActionType<Response> action,
1649 final Request request, final ActionListener<Response> listener) {
1650 client.execute(action, request, listener);
1651 }
1652
1653 @Override
1654 public FieldCapabilitiesRequestBuilder prepareFieldCaps(final String... indices) {
1655 return client.prepareFieldCaps(indices);
1656 }
1657
1658 @Override
1659 public ActionFuture<FieldCapabilitiesResponse> fieldCaps(final FieldCapabilitiesRequest request) {
1660 return client.fieldCaps(request);
1661 }
1662
1663 @Override
1664 public void fieldCaps(final FieldCapabilitiesRequest request, final ActionListener<FieldCapabilitiesResponse> listener) {
1665 client.fieldCaps(request, listener);
1666 }
1667
1668 @Override
1669 public BulkRequestBuilder prepareBulk(final String globalIndex, final String globalType) {
1670 return client.prepareBulk(globalIndex, globalType);
1671 }
1672 }