1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16 package org.codelibs.fess.es.config.allcommon;
17
18 import java.time.LocalDateTime;
19 import java.util.ArrayList;
20 import java.util.Collection;
21 import java.util.Date;
22 import java.util.Iterator;
23 import java.util.List;
24 import java.util.ListIterator;
25 import java.util.Map;
26 import java.util.function.Function;
27
28 import javax.annotation.Resource;
29
30 import org.codelibs.fess.es.config.allcommon.EsAbstractEntity.DocMeta;
31 import org.codelibs.fess.es.config.allcommon.EsAbstractEntity.RequestOptionCall;
32 import org.dbflute.Entity;
33 import org.dbflute.bhv.AbstractBehaviorWritable;
34 import org.dbflute.bhv.readable.EntityRowHandler;
35 import org.dbflute.bhv.writable.DeleteOption;
36 import org.dbflute.bhv.writable.InsertOption;
37 import org.dbflute.bhv.writable.UpdateOption;
38 import org.dbflute.cbean.ConditionBean;
39 import org.dbflute.cbean.coption.CursorSelectOption;
40 import org.dbflute.cbean.result.ListResultBean;
41 import org.dbflute.exception.FetchingOverSafetySizeException;
42 import org.dbflute.exception.IllegalBehaviorStateException;
43 import org.dbflute.util.DfTypeUtil;
44 import org.elasticsearch.action.DocWriteResponse.Result;
45 import org.elasticsearch.action.admin.indices.refresh.RefreshResponse;
46 import org.elasticsearch.action.bulk.BulkItemResponse;
47 import org.elasticsearch.action.bulk.BulkRequestBuilder;
48 import org.elasticsearch.action.bulk.BulkResponse;
49 import org.elasticsearch.action.delete.DeleteRequestBuilder;
50 import org.elasticsearch.action.delete.DeleteResponse;
51 import org.elasticsearch.action.index.IndexRequestBuilder;
52 import org.elasticsearch.action.index.IndexResponse;
53 import org.elasticsearch.action.search.SearchRequestBuilder;
54 import org.elasticsearch.action.search.SearchResponse;
55 import org.elasticsearch.action.update.UpdateRequestBuilder;
56 import org.elasticsearch.client.Client;
57 import org.elasticsearch.search.SearchHit;
58 import org.elasticsearch.search.SearchHits;
59
60
61
62
63
64
65 public abstract class EsAbstractBehavior<ENTITY extends Entity, CB extends ConditionBean> extends AbstractBehaviorWritable<ENTITY, CB> {
66
67 @Resource
68 private Client client;
69
70 protected int sizeForDelete = 100;
71 protected String scrollForDelete = "1m";
72 protected int sizeForCursor = 100;
73 protected String scrollForCursor = "1m";
74 protected String searchTimeout = "3m";
75 protected String indexTimeout = "3m";
76 protected String scrollSearchTimeout = "3m";
77 protected String bulkTimeout = "3m";
78 protected String deleteTimeout = "3m";
79 protected String refreshTimeout = "1m";
80
81 protected abstract String asEsIndex();
82
83 protected abstract String asEsIndexType();
84
85 protected abstract String asEsSearchType();
86
87 protected abstract <RESULT extends ENTITY> RESULT createEntity(Map<String, Object> source, Class<? extends RESULT> entityType);
88
89
90
91
92 public RefreshResponse refresh() {
93 return client.admin().indices().prepareRefresh(asEsIndex()).execute().actionGet(refreshTimeout);
94 }
95
96
97
98
99 @Override
100 protected int delegateSelectCountUniquely(final ConditionBean cb) {
101
102 final SearchRequestBuilder builder = client.prepareSearch(asEsIndex()).setTypes(asEsSearchType());
103 final EsAbstractConditionBean esCb = (EsAbstractConditionBean) cb;
104 if (esCb.getPreference() != null) {
105 builder.setPreference(esCb.getPreference());
106 }
107 return (int) esCb.build(builder).execute().actionGet(searchTimeout).getHits().getTotalHits();
108 }
109
110 @Override
111 protected <RESULT extends ENTITY> RESULT delegateSelectEntity(final ConditionBean cb, final Class<? extends RESULT> entityType) {
112 final List<? extends RESULT> list = delegateSelectList(cb, entityType);
113 if (list.isEmpty()) {
114 return null;
115 }
116 if (list.size() >= 2) {
117 String msg = "The size of selected list is over 1: " + list.size();
118 throw new FetchingOverSafetySizeException(msg, 1);
119 }
120 return list.get(0);
121 }
122
123 @Override
124 protected <RESULT extends ENTITY> List<RESULT> delegateSelectList(final ConditionBean cb, final Class<? extends RESULT> entityType) {
125
126 final SearchRequestBuilder builder = client.prepareSearch(asEsIndex()).setTypes(asEsSearchType());
127 final int from;
128 final int size;
129 if (cb.isFetchScopeEffective()) {
130 from = cb.getPageStartIndex();
131 size = cb.getFetchSize();
132 } else {
133 from = 0;
134 size = 10;
135 }
136 builder.setFrom(from);
137 builder.setSize(size);
138 final EsAbstractConditionBean esCb = (EsAbstractConditionBean) cb;
139 if (esCb.getPreference() != null) {
140 builder.setPreference(esCb.getPreference());
141 }
142 esCb.request().build(builder);
143 final SearchResponse response = esCb.build(builder).execute().actionGet(searchTimeout);
144
145 final EsPagingResultBean<RESULT> list = new EsPagingResultBean<>(builder);
146 final SearchHits searchHits = response.getHits();
147 searchHits.forEach(hit -> {
148 final Map<String, Object> source = hit.getSource();
149 final RESULT entity = createEntity(source, entityType);
150 final DocMeta docMeta = ((EsAbstractEntity) entity).asDocMeta();
151 docMeta.id(hit.getId());
152 docMeta.version(hit.getVersion());
153 list.add(entity);
154 });
155
156 list.setPageSize(size);
157 list.setAllRecordCount((int) searchHits.getTotalHits());
158 list.setCurrentPageNumber(cb.getFetchPageNumber());
159
160 list.setTook(response.getTookInMillis());
161 list.setTotalShards(response.getTotalShards());
162 list.setSuccessfulShards(response.getSuccessfulShards());
163 list.setFailedShards(response.getFailedShards());
164
165 list.setAggregation(response.getAggregations());
166
167
168
169 return list;
170 }
171
172 @Override
173 protected <RESULT extends ENTITY> void helpSelectCursorHandlingByPaging(CB cb, EntityRowHandler<RESULT> handler,
174 Class<? extends RESULT> entityType, CursorSelectOption option) {
175 delegateSelectCursor(cb, handler, entityType);
176 }
177
178 @Override
179 protected <RESULT extends ENTITY> void delegateSelectCursor(final ConditionBean cb, final EntityRowHandler<RESULT> handler,
180 final Class<? extends RESULT> entityType) {
181 delegateBulkRequest(cb, searchHits -> {
182 searchHits.forEach(hit -> {
183 if (handler.isBreakCursor()) {
184 return;
185 }
186 final Map<String, Object> source = hit.getSource();
187 final RESULT entity = createEntity(source, entityType);
188 final DocMeta docMeta = ((EsAbstractEntity) entity).asDocMeta();
189 docMeta.id(hit.getId());
190 docMeta.version(hit.getVersion());
191 handler.handle(entity);
192 });
193
194 return !handler.isBreakCursor();
195 });
196 }
197
198 protected <RESULT extends ENTITY> void delegateSelectBulk(final ConditionBean cb, final EntityRowHandler<List<RESULT>> handler,
199 final Class<? extends RESULT> entityType) {
200 assertCBStateValid(cb);
201 assertObjectNotNull("entityRowHandler", handler);
202 assertSpecifyDerivedReferrerEntityProperty(cb, entityType);
203 assertObjectNotNull("entityRowHandler", handler);
204 delegateBulkRequest(cb, searchHits -> {
205 List<RESULT> list = new ArrayList<>();
206 searchHits.forEach(hit -> {
207 final Map<String, Object> source = hit.getSource();
208 final RESULT entity = createEntity(source, entityType);
209 final DocMeta docMeta = ((EsAbstractEntity) entity).asDocMeta();
210 docMeta.id(hit.getId());
211 docMeta.version(hit.getVersion());
212 list.add(entity);
213 });
214
215 handler.handle(list);
216 return !handler.isBreakCursor();
217 });
218 }
219
220 protected void delegateBulkRequest(final ConditionBean cb, Function<SearchHits, Boolean> handler) {
221 SearchResponse response = null;
222 while (true) {
223 if (response == null) {
224 final SearchRequestBuilder builder =
225 client.prepareSearch(asEsIndex()).setTypes(asEsIndexType()).setScroll(scrollForCursor).setSize(sizeForCursor);
226 final EsAbstractConditionBean esCb = (EsAbstractConditionBean) cb;
227 if (esCb.getPreference() != null) {
228 builder.setPreference(esCb.getPreference());
229 }
230 esCb.request().build(builder);
231 response = esCb.build(builder).execute().actionGet(scrollSearchTimeout);
232 } else {
233 final String scrollId = response.getScrollId();
234 response = client.prepareSearchScroll(scrollId).setScroll(scrollForDelete).execute().actionGet(scrollSearchTimeout);
235 }
236 final SearchHits searchHits = response.getHits();
237 final SearchHit[] hits = searchHits.getHits();
238 if (hits.length == 0) {
239 break;
240 }
241
242 if (!handler.apply(searchHits)) {
243 break;
244 }
245 }
246 }
247
248 @Override
249 protected Number doReadNextVal() {
250 final String msg = "This table is NOT related to sequence: " + asEsIndexType();
251 throw new UnsupportedOperationException(msg);
252 }
253
254 @Override
255 protected <RESULT extends Entity> ListResultBean<RESULT> createListResultBean(final ConditionBean cb, final List<RESULT> selectedList) {
256 if (selectedList instanceof EsPagingResultBean) {
257 return (ListResultBean<RESULT>) selectedList;
258 }
259 throw new IllegalBehaviorStateException("selectedList is not EsPagingResultBean.");
260 }
261
262
263
264
265 @Override
266 protected int delegateInsert(final Entity entity, final InsertOption<? extends ConditionBean> option) {
267 final EsAbstractEntity esEntity = (EsAbstractEntity) entity;
268 IndexRequestBuilder builder = createInsertRequest(esEntity);
269
270 final IndexResponse response = builder.execute().actionGet(indexTimeout);
271 esEntity.asDocMeta().id(response.getId());
272 return response.getResult() == Result.CREATED ? 1 : 0;
273 }
274
275 protected IndexRequestBuilder createInsertRequest(final EsAbstractEntity esEntity) {
276 final IndexRequestBuilder builder = client.prepareIndex(asEsIndex(), asEsIndexType()).setSource(toSource(esEntity));
277 final String id = esEntity.asDocMeta().id();
278 if (id != null) {
279 builder.setId(id);
280 }
281 final RequestOptionCall<IndexRequestBuilder> indexOption = esEntity.asDocMeta().indexOption();
282 if (indexOption != null) {
283 indexOption.callback(builder);
284 }
285 return builder;
286 }
287
288 @Override
289 protected int delegateUpdate(final Entity entity, final UpdateOption<? extends ConditionBean> option) {
290 final EsAbstractEntity esEntity = (EsAbstractEntity) entity;
291 final IndexRequestBuilder builder = createUpdateRequest(esEntity);
292
293 final IndexResponse response = builder.execute().actionGet(indexTimeout);
294 long version = response.getVersion();
295 if (version != -1) {
296 esEntity.asDocMeta().version(version);
297 }
298 return 1;
299 }
300
301 protected IndexRequestBuilder createUpdateRequest(final EsAbstractEntity esEntity) {
302 final IndexRequestBuilder builder =
303 client.prepareIndex(asEsIndex(), asEsIndexType(), esEntity.asDocMeta().id()).setSource(toSource(esEntity));
304 final RequestOptionCall<IndexRequestBuilder> indexOption = esEntity.asDocMeta().indexOption();
305 if (indexOption != null) {
306 indexOption.callback(builder);
307 }
308 final Long version = esEntity.asDocMeta().version();
309 if (version != null && version.longValue() != -1) {
310 builder.setVersion(version);
311 }
312 return builder;
313 }
314
315 protected Map<String, Object> toSource(final EsAbstractEntity esEntity) {
316 return esEntity.toSource();
317 }
318
319 @Override
320 protected int delegateDelete(final Entity entity, final DeleteOption<? extends ConditionBean> option) {
321 final EsAbstractEntity esEntity = (EsAbstractEntity) entity;
322 final DeleteRequestBuilder builder = createDeleteRequest(esEntity);
323
324 final DeleteResponse response = builder.execute().actionGet(deleteTimeout);
325 return response.getResult() == Result.DELETED ? 1 : 0;
326 }
327
328 protected DeleteRequestBuilder createDeleteRequest(final EsAbstractEntity esEntity) {
329 final DeleteRequestBuilder builder = client.prepareDelete(asEsIndex(), asEsIndexType(), esEntity.asDocMeta().id());
330 final RequestOptionCall<DeleteRequestBuilder> deleteOption = esEntity.asDocMeta().deleteOption();
331 if (deleteOption != null) {
332 deleteOption.callback(builder);
333 }
334 return builder;
335 }
336
337 @Override
338 protected int delegateQueryDelete(final ConditionBean cb, final DeleteOption<? extends ConditionBean> option) {
339 SearchResponse response = null;
340 int count = 0;
341 while (true) {
342 if (response == null) {
343 final SearchRequestBuilder builder =
344 client.prepareSearch(asEsIndex()).setTypes(asEsIndexType()).setScroll(scrollForDelete).setSize(sizeForDelete);
345 final EsAbstractConditionBean esCb = (EsAbstractConditionBean) cb;
346 if (esCb.getPreference() != null) {
347 esCb.setPreference(esCb.getPreference());
348 }
349 esCb.request().build(builder);
350 response = esCb.build(builder).execute().actionGet(scrollSearchTimeout);
351 } else {
352 final String scrollId = response.getScrollId();
353 response = client.prepareSearchScroll(scrollId).setScroll(scrollForDelete).execute().actionGet(scrollSearchTimeout);
354 }
355 final SearchHits searchHits = response.getHits();
356 final SearchHit[] hits = searchHits.getHits();
357 if (hits.length == 0) {
358 break;
359 }
360
361 final BulkRequestBuilder bulkRequest = client.prepareBulk();
362 for (final SearchHit hit : hits) {
363 bulkRequest.add(client.prepareDelete(asEsIndex(), asEsIndexType(), hit.getId()));
364 }
365 count += hits.length;
366 final BulkResponse bulkResponse = bulkRequest.execute().actionGet(bulkTimeout);
367 if (bulkResponse.hasFailures()) {
368 throw new IllegalBehaviorStateException(bulkResponse.buildFailureMessage());
369 }
370 }
371 return count;
372 }
373
374 protected int[] delegateBatchInsert(final List<? extends Entity> entityList, final InsertOption<? extends ConditionBean> option) {
375 if (entityList.isEmpty()) {
376 return new int[] {};
377 }
378 return delegateBatchRequest(entityList, esEntity -> {
379 return createInsertRequest(esEntity);
380 });
381 }
382
383 protected int[] delegateBatchUpdate(List<? extends Entity> entityList, UpdateOption<? extends ConditionBean> option) {
384 if (entityList.isEmpty()) {
385 return new int[] {};
386 }
387 return delegateBatchRequest(entityList, esEntity -> {
388 return createUpdateRequest(esEntity);
389 });
390 }
391
392 protected int[] delegateBatchDelete(List<? extends Entity> entityList, DeleteOption<? extends ConditionBean> option) {
393 if (entityList.isEmpty()) {
394 return new int[] {};
395 }
396 return delegateBatchRequest(entityList, esEntity -> {
397 return createDeleteRequest(esEntity);
398 });
399 }
400
401 protected <BUILDER> int[] delegateBatchRequest(final List<? extends Entity> entityList, Function<EsAbstractEntity, BUILDER> call) {
402 @SuppressWarnings("unchecked")
403 final BulkList<? extends Entity, BUILDER> bulkList = (BulkList<? extends Entity, BUILDER>) entityList;
404 final RequestOptionCall<BUILDER> builderEntityCall = bulkList.getEntityCall();
405 final BulkRequestBuilder bulkBuilder = client.prepareBulk();
406 for (final Entity entity : entityList) {
407 final EsAbstractEntity esEntity = (EsAbstractEntity) entity;
408 BUILDER builder = call.apply(esEntity);
409 if (builder instanceof IndexRequestBuilder) {
410 if (builderEntityCall != null) {
411 builderEntityCall.callback(builder);
412 }
413 bulkBuilder.add((IndexRequestBuilder) builder);
414 } else if (builder instanceof UpdateRequestBuilder) {
415 if (builderEntityCall != null) {
416 builderEntityCall.callback(builder);
417 }
418 bulkBuilder.add((UpdateRequestBuilder) builder);
419 } else if (builder instanceof DeleteRequestBuilder) {
420 if (builderEntityCall != null) {
421 builderEntityCall.callback(builder);
422 }
423 bulkBuilder.add((DeleteRequestBuilder) builder);
424 }
425 }
426 final RequestOptionCall<BulkRequestBuilder> builderCall = bulkList.getCall();
427 if (builderCall != null) {
428 builderCall.callback(bulkBuilder);
429 }
430
431 final BulkResponse response = bulkBuilder.execute().actionGet(bulkTimeout);
432 final BulkItemResponse[] itemResponses = response.getItems();
433 if (itemResponses.length != entityList.size()) {
434 throw new IllegalStateException("Invalid response size: " + itemResponses.length + " != " + entityList.size());
435 }
436 final int[] results = new int[itemResponses.length];
437 for (int i = 0; i < itemResponses.length; i++) {
438 final BulkItemResponse itemResponse = itemResponses[i];
439 final Entity entity = entityList.get(i);
440 if (entity instanceof EsAbstractEntity) {
441 ((EsAbstractEntity) entity).asDocMeta().id(itemResponse.getId());
442 }
443 results[i] = itemResponse.isFailed() ? 0 : 1;
444 }
445 return results;
446 }
447
448
449 @Override
450 protected UpdateOption<CB> createPlainUpdateOption() {
451 UpdateOption<CB> updateOption = new UpdateOption<CB>();
452 updateOption.xtoBeCompatibleBatchUpdateDefaultEveryColumn();
453 return updateOption;
454 }
455
456 protected boolean isCompatibleBatchInsertDefaultEveryColumn() {
457 return true;
458 }
459
460 public void setSizeForDelete(int sizeForDelete) {
461 this.sizeForDelete = sizeForDelete;
462 }
463
464 public void setScrollForDelete(String scrollForDelete) {
465 this.scrollForDelete = scrollForDelete;
466 }
467
468 public void setSizeForCursor(int sizeForCursor) {
469 this.sizeForCursor = sizeForCursor;
470 }
471
472 public void setScrollForCursor(String scrollForCursor) {
473 this.scrollForCursor = scrollForCursor;
474 }
475
476 public void setSearchTimeout(String searchTimeout) {
477 this.searchTimeout = searchTimeout;
478 }
479
480 public void setIndexTimeout(String indexTimeout) {
481 this.indexTimeout = indexTimeout;
482 }
483
484 public void setScrollSearchTimeout(String scrollSearchTimeout) {
485 this.scrollSearchTimeout = scrollSearchTimeout;
486 }
487
488 public void setBulkTimeout(String bulkTimeout) {
489 this.bulkTimeout = bulkTimeout;
490 }
491
492 public void setDeleteTimeout(String deleteTimeout) {
493 this.deleteTimeout = deleteTimeout;
494 }
495
496 public void setRefreshTimeout(String refreshTimeout) {
497 this.refreshTimeout = refreshTimeout;
498 }
499
500
501
502
503 protected String[] toStringArray(final Object value) {
504 if (value instanceof String[]) {
505 return (String[]) value;
506 } else if (value instanceof List) {
507 return ((List<?>) value).stream().map(v -> v.toString()).toArray(n -> new String[n]);
508 }
509 String str = DfTypeUtil.toString(value);
510 if (str == null) {
511 return null;
512 }
513 return new String[] { str };
514 }
515
516 protected LocalDateTime toLocalDateTime(Object value) {
517 return DfTypeUtil.toLocalDateTime(value);
518 }
519
520 protected Date toDate(Object value) {
521 return DfTypeUtil.toDate(value);
522 }
523
524 public static class BulkList<E, B> implements List<E> {
525
526 private final List<E> parent;
527
528 private final RequestOptionCall<BulkRequestBuilder> call;
529
530 private final RequestOptionCall<B> entityCall;
531
532 public BulkList(final List<E> parent, final RequestOptionCall<BulkRequestBuilder> call, final RequestOptionCall<B> entityCall) {
533 this.parent = parent;
534 this.entityCall = entityCall;
535 this.call = call;
536 }
537
538 public int size() {
539 return parent.size();
540 }
541
542 public boolean isEmpty() {
543 return parent.isEmpty();
544 }
545
546 public boolean contains(final Object o) {
547 return parent.contains(o);
548 }
549
550 public Iterator<E> iterator() {
551 return parent.iterator();
552 }
553
554 public Object[] toArray() {
555 return parent.toArray();
556 }
557
558 public <T> T[] toArray(final T[] a) {
559 return parent.toArray(a);
560 }
561
562 public boolean add(final E e) {
563 return parent.add(e);
564 }
565
566 public boolean remove(final Object o) {
567 return parent.remove(o);
568 }
569
570 public boolean containsAll(final Collection<?> c) {
571 return parent.containsAll(c);
572 }
573
574 public boolean addAll(final Collection<? extends E> c) {
575 return parent.addAll(c);
576 }
577
578 public boolean addAll(final int index, final Collection<? extends E> c) {
579 return parent.addAll(index, c);
580 }
581
582 public boolean removeAll(final Collection<?> c) {
583 return parent.removeAll(c);
584 }
585
586 public boolean retainAll(final Collection<?> c) {
587 return parent.retainAll(c);
588 }
589
590 public void clear() {
591 parent.clear();
592 }
593
594 public boolean equals(final Object o) {
595 return parent.equals(o);
596 }
597
598 public int hashCode() {
599 return parent.hashCode();
600 }
601
602 public E get(final int index) {
603 return parent.get(index);
604 }
605
606 public E set(final int index, final E element) {
607 return parent.set(index, element);
608 }
609
610 public void add(final int index, final E element) {
611 parent.add(index, element);
612 }
613
614 public E remove(final int index) {
615 return parent.remove(index);
616 }
617
618 public int indexOf(final Object o) {
619 return parent.indexOf(o);
620 }
621
622 public int lastIndexOf(final Object o) {
623 return parent.lastIndexOf(o);
624 }
625
626 public ListIterator<E> listIterator() {
627 return parent.listIterator();
628 }
629
630 public ListIterator<E> listIterator(final int index) {
631 return parent.listIterator(index);
632 }
633
634 public List<E> subList(final int fromIndex, final int toIndex) {
635 return parent.subList(fromIndex, toIndex);
636 }
637
638 public RequestOptionCall<BulkRequestBuilder> getCall() {
639 return call;
640 }
641
642 public RequestOptionCall<B> getEntityCall() {
643 return entityCall;
644 }
645 }
646 }