View Javadoc
1   /*
2    * Copyright 2012-2021 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.helper;
17  
18  import java.util.ArrayList;
19  import java.util.Collections;
20  import java.util.HashMap;
21  import java.util.List;
22  import java.util.Map;
23  
24  import org.apache.logging.log4j.LogManager;
25  import org.apache.logging.log4j.Logger;
26  import org.codelibs.core.lang.StringUtil;
27  import org.codelibs.core.lang.ThreadUtil;
28  import org.codelibs.fesen.index.query.QueryBuilder;
29  import org.codelibs.fesen.index.query.QueryBuilders;
30  import org.codelibs.fess.Constants;
31  import org.codelibs.fess.app.service.FailureUrlService;
32  import org.codelibs.fess.ds.DataStore;
33  import org.codelibs.fess.ds.DataStoreFactory;
34  import org.codelibs.fess.ds.callback.IndexUpdateCallback;
35  import org.codelibs.fess.es.client.SearchEngineClient;
36  import org.codelibs.fess.es.config.exentity.DataConfig;
37  import org.codelibs.fess.mylasta.direction.FessConfig;
38  import org.codelibs.fess.util.ComponentUtil;
39  
40  public class DataIndexHelper {
41  
42      private static final Logger logger = LogManager.getLogger(DataIndexHelper.class);
43  
44      private static final String DELETE_OLD_DOCS = "delete_old_docs";
45  
46      protected long crawlingExecutionInterval = Constants.DEFAULT_CRAWLING_EXECUTION_INTERVAL;
47  
48      protected int crawlerPriority = Thread.NORM_PRIORITY;
49  
50      protected final List<DataCrawlingThread> dataCrawlingThreadList = Collections.synchronizedList(new ArrayList<DataCrawlingThread>());
51  
52      public void crawl(final String sessionId) {
53          final List<DataConfig> configList = ComponentUtil.getCrawlingConfigHelper().getAllDataConfigList();
54  
55          if (configList.isEmpty()) {
56              // nothing
57              if (logger.isInfoEnabled()) {
58                  logger.info("No crawling target data.");
59              }
60              return;
61          }
62  
63          doCrawl(sessionId, configList);
64      }
65  
66      public void crawl(final String sessionId, final List<String> configIdList) {
67          final List<DataConfig> configList = ComponentUtil.getCrawlingConfigHelper().getDataConfigListByIds(configIdList);
68  
69          if (configList.isEmpty()) {
70              // nothing
71              if (logger.isInfoEnabled()) {
72                  logger.info("No crawling target urls.");
73              }
74              return;
75          }
76  
77          doCrawl(sessionId, configList);
78      }
79  
80      protected void doCrawl(final String sessionId, final List<DataConfig> configList) {
81          final int multiprocessCrawlingCount = ComponentUtil.getFessConfig().getCrawlingThreadCount();
82  
83          final long startTime = System.currentTimeMillis();
84  
85          final IndexUpdateCallback indexUpdateCallback = ComponentUtil.getComponent(IndexUpdateCallback.class);
86  
87          final List<String> sessionIdList = new ArrayList<>();
88          dataCrawlingThreadList.clear();
89          final List<String> dataCrawlingThreadStatusList = new ArrayList<>();
90          for (final DataConfig dataConfig : configList) {
91              final Map<String, String> initParamMap = new HashMap<>();
92              final String sid = ComponentUtil.getCrawlingConfigHelper().store(sessionId, dataConfig);
93              sessionIdList.add(sid);
94  
95              initParamMap.put(Constants.SESSION_ID, sessionId);
96              initParamMap.put(Constants.CRAWLING_INFO_ID, sid);
97  
98              final DataCrawlingThread dataCrawlingThread = new DataCrawlingThread(dataConfig, indexUpdateCallback, initParamMap);
99              dataCrawlingThread.setPriority(crawlerPriority);
100             dataCrawlingThread.setName(sid);
101             dataCrawlingThread.setDaemon(true);
102 
103             dataCrawlingThreadList.add(dataCrawlingThread);
104             dataCrawlingThreadStatusList.add(Constants.READY);
105 
106         }
107 
108         final SystemHelper systemHelper = ComponentUtil.getSystemHelper();
109 
110         int startedCrawlerNum = 0;
111         int activeCrawlerNum = 0;
112         while (startedCrawlerNum < dataCrawlingThreadList.size()) {
113             // Force to stop crawl
114             if (systemHelper.isForceStop()) {
115                 for (final DataCrawlingThread crawlerThread : dataCrawlingThreadList) {
116                     crawlerThread.stopCrawling();
117                 }
118                 break;
119             }
120 
121             if (activeCrawlerNum < multiprocessCrawlingCount) {
122                 // start crawling
123                 dataCrawlingThreadList.get(startedCrawlerNum).start();
124                 dataCrawlingThreadStatusList.set(startedCrawlerNum, Constants.RUNNING);
125                 startedCrawlerNum++;
126                 activeCrawlerNum++;
127                 ThreadUtil.sleep(crawlingExecutionInterval);
128                 continue;
129             }
130 
131             // check status
132             for (int i = 0; i < startedCrawlerNum; i++) {
133                 if (!dataCrawlingThreadList.get(i).isRunning() && Constants.RUNNING.equals(dataCrawlingThreadStatusList.get(i))) {
134                     dataCrawlingThreadList.get(i).awaitTermination();
135                     dataCrawlingThreadStatusList.set(i, Constants.DONE);
136                     activeCrawlerNum--;
137                 }
138             }
139             ThreadUtil.sleep(crawlingExecutionInterval);
140         }
141 
142         boolean finishedAll = false;
143         while (!finishedAll) {
144             finishedAll = true;
145             for (int i = 0; i < dataCrawlingThreadList.size(); i++) {
146                 dataCrawlingThreadList.get(i).awaitTermination(crawlingExecutionInterval);
147                 if (!dataCrawlingThreadList.get(i).isRunning() && Constants.RUNNING.equals(dataCrawlingThreadStatusList.get(i))) {
148                     dataCrawlingThreadStatusList.set(i, Constants.DONE);
149                 }
150                 if (!Constants.DONE.equals(dataCrawlingThreadStatusList.get(i))) {
151                     finishedAll = false;
152                 }
153             }
154         }
155         dataCrawlingThreadList.clear();
156         dataCrawlingThreadStatusList.clear();
157 
158         // put cralwing info
159         final CrawlingInfoHelper crawlingInfoHelper = ComponentUtil.getCrawlingInfoHelper();
160 
161         final long execTime = System.currentTimeMillis() - startTime;
162         crawlingInfoHelper.putToInfoMap(Constants.DATA_CRAWLING_EXEC_TIME, Long.toString(execTime));
163         if (logger.isInfoEnabled()) {
164             logger.info("[EXEC TIME] crawling time: {}ms", execTime);
165         }
166 
167         crawlingInfoHelper.putToInfoMap(Constants.DATA_INDEX_EXEC_TIME, Long.toString(indexUpdateCallback.getExecuteTime()));
168         crawlingInfoHelper.putToInfoMap(Constants.DATA_INDEX_SIZE, Long.toString(indexUpdateCallback.getDocumentSize()));
169 
170         for (final String sid : sessionIdList) {
171             // remove config
172             ComponentUtil.getCrawlingConfigHelper().remove(sid);
173         }
174 
175     }
176 
177     protected static class DataCrawlingThread extends Thread {
178 
179         private final DataConfig dataConfig;
180 
181         private final IndexUpdateCallback indexUpdateCallback;
182 
183         private final Map<String, String> initParamMap;
184 
185         protected boolean finished = false;
186 
187         protected boolean running = false;
188 
189         private DataStore dataStore;
190 
191         protected DataCrawlingThread(final DataConfig dataConfig, final IndexUpdateCallback indexUpdateCallback,
192                 final Map<String, String> initParamMap) {
193             this.dataConfig = dataConfig;
194             this.indexUpdateCallback = indexUpdateCallback;
195             this.initParamMap = initParamMap;
196         }
197 
198         @Override
199         public void run() {
200             running = true;
201             try {
202                 process();
203             } finally {
204                 running = false;
205                 finished = true;
206             }
207         }
208 
209         protected void process() {
210             final DataStoreFactory dataStoreFactory = ComponentUtil.getDataStoreFactory();
211             dataStore = dataStoreFactory.getDataStore(dataConfig.getHandlerName());
212             if (dataStore == null) {
213                 logger.error("DataStore({}) is not found.", dataConfig.getHandlerName());
214             } else {
215                 try {
216                     dataStore.store(dataConfig, indexUpdateCallback, initParamMap);
217                 } catch (final Throwable e) {
218                     logger.error("Failed to process a data crawling: {}", dataConfig.getName(), e);
219                     ComponentUtil.getComponent(FailureUrlService.class).store(dataConfig, e.getClass().getCanonicalName(),
220                             dataConfig.getConfigId() + ":" + dataConfig.getName(), e);
221                 } finally {
222                     indexUpdateCallback.commit();
223                     deleteOldDocs();
224                 }
225             }
226         }
227 
228         private void deleteOldDocs() {
229             if (Constants.FALSE.equals(initParamMap.get(DELETE_OLD_DOCS))) {
230                 return;
231             }
232             final String sessionId = initParamMap.get(Constants.SESSION_ID);
233             if (StringUtil.isBlank(sessionId)) {
234                 logger.warn("Invalid sessionId at {}", dataConfig);
235                 return;
236             }
237             final FessConfig fessConfig = ComponentUtil.getFessConfig();
238             final QueryBuilder queryBuilder =
239                     QueryBuilders.boolQuery().must(QueryBuilders.termQuery(fessConfig.getIndexFieldConfigId(), dataConfig.getConfigId()))
240                             .must(QueryBuilders.boolQuery().mustNot(QueryBuilders.rangeQuery(fessConfig.getIndexFieldExpires()).gt("now"))
241                                     .mustNot(QueryBuilders.existsQuery(fessConfig.getIndexFieldExpires())))
242                             .mustNot(QueryBuilders.termQuery(fessConfig.getIndexFieldSegment(), sessionId));
243             try {
244                 final SearchEngineClient searchEngineClient = ComponentUtil.getSearchEngineClient();
245                 final String index = fessConfig.getIndexDocumentUpdateIndex();
246                 searchEngineClient.admin().indices().prepareRefresh(index).execute().actionGet();
247                 final long numOfDeleted = searchEngineClient.deleteByQuery(index, queryBuilder);
248                 logger.info("Deleted {} old docs.", numOfDeleted);
249             } catch (final Exception e) {
250                 logger.error("Could not delete old docs at {}", dataConfig, e);
251             }
252         }
253 
254         public boolean isFinished() {
255             return finished;
256         }
257 
258         public void stopCrawling() {
259             if (dataStore != null) {
260                 dataStore.stop();
261             }
262         }
263 
264         public String getCrawlingInfoId() {
265             return initParamMap.get(Constants.CRAWLING_INFO_ID);
266         }
267 
268         public boolean isRunning() {
269             return running;
270         }
271 
272         public void awaitTermination() {
273             try {
274                 join();
275             } catch (final InterruptedException e) {
276                 if (logger.isDebugEnabled()) {
277                     logger.debug("Interrupted.", e);
278                 }
279             }
280         }
281 
282         public void awaitTermination(final long mills) {
283             try {
284                 join(mills);
285             } catch (final InterruptedException e) {
286                 if (logger.isDebugEnabled()) {
287                     logger.debug("Interrupted.", e);
288                 }
289             }
290         }
291     }
292 
293     public void setCrawlingExecutionInterval(final long crawlingExecutionInterval) {
294         this.crawlingExecutionInterval = crawlingExecutionInterval;
295     }
296 
297     public void setCrawlerPriority(final int crawlerPriority) {
298         this.crawlerPriority = crawlerPriority;
299     }
300 }