View Javadoc
1   /*
2    * Copyright 2012-2017 CodeLibs Project and the Others.
3    *
4    * Licensed under the Apache License, Version 2.0 (the "License");
5    * you may not use this file except in compliance with the License.
6    * You may obtain a copy of the License at
7    *
8    *     http://www.apache.org/licenses/LICENSE-2.0
9    *
10   * Unless required by applicable law or agreed to in writing, software
11   * distributed under the License is distributed on an "AS IS" BASIS,
12   * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND,
13   * either express or implied. See the License for the specific language
14   * governing permissions and limitations under the License.
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   * @param <ENTITY> The type of entity.
62   * @param <CB> The type of condition-bean.
63   * @author ESFlute (using FreeGen)
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      //                                                                       Elasticsearch
91      //                                                                              ======
92      public RefreshResponse refresh() {
93          return client.admin().indices().prepareRefresh(asEsIndex()).execute().actionGet(refreshTimeout);
94      }
95  
96      // ===================================================================================
97      //                                                                              Select
98      //                                                                              ======
99      @Override
100     protected int delegateSelectCountUniquely(final ConditionBean cb) {
101         // #pending check response and cast problem
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); // immediatly catched by caller and tranlated 
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         // #pending check response
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         // #pending others
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     //                                                                              Update
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     // to suppress xacceptUpdateColumnModifiedPropertiesIfNeeds()'s specify process
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     //                                                                        Assist Logic
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 }