1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16 package org.codelibs.fess.helper;
17
18 import java.util.ArrayList;
19 import java.util.Collections;
20 import java.util.List;
21
22 import org.apache.logging.log4j.LogManager;
23 import org.apache.logging.log4j.Logger;
24 import org.codelibs.core.lang.StringUtil;
25 import org.codelibs.core.lang.ThreadUtil;
26 import org.codelibs.fess.Constants;
27 import org.codelibs.fess.app.service.FailureUrlService;
28 import org.codelibs.fess.ds.DataStore;
29 import org.codelibs.fess.ds.DataStoreFactory;
30 import org.codelibs.fess.ds.callback.IndexUpdateCallback;
31 import org.codelibs.fess.entity.DataStoreParams;
32 import org.codelibs.fess.mylasta.direction.FessConfig;
33 import org.codelibs.fess.opensearch.client.SearchEngineClient;
34 import org.codelibs.fess.opensearch.config.exentity.DataConfig;
35 import org.codelibs.fess.util.ComponentUtil;
36 import org.opensearch.index.query.BoolQueryBuilder;
37 import org.opensearch.index.query.QueryBuilder;
38 import org.opensearch.index.query.QueryBuilders;
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55 public class DataIndexHelper {
56
57
58 private static final Logger logger = LogManager.getLogger(DataIndexHelper.class);
59
60
61 private static final String DELETE_OLD_DOCS = "delete_old_docs";
62
63
64 private static final String KEEP_EXPIRES_DOCS = "keep_expires_docs";
65
66
67
68
69
70
71 protected long crawlingExecutionInterval = Constants.DEFAULT_CRAWLING_EXECUTION_INTERVAL;
72
73
74
75
76
77 protected int crawlerPriority = Thread.NORM_PRIORITY;
78
79
80
81
82
83 protected final List<DataCrawlingThread> dataCrawlingThreadList = Collections.synchronizedList(new ArrayList<>());
84
85
86
87
88
89
90 public DataIndexHelper() {
91
92 }
93
94
95
96
97
98
99
100
101 public void crawl(final String sessionId) {
102 final List<DataConfig> configList = ComponentUtil.getCrawlingConfigHelper().getAllDataConfigList();
103
104 if (configList.isEmpty()) {
105
106 if (logger.isInfoEnabled()) {
107 logger.info("No crawling target data.");
108 }
109 return;
110 }
111
112 doCrawl(sessionId, configList);
113 }
114
115
116
117
118
119
120
121
122
123 public void crawl(final String sessionId, final List<String> configIdList) {
124 final List<DataConfig> configList = ComponentUtil.getCrawlingConfigHelper().getDataConfigListByIds(configIdList);
125
126 if (configList.isEmpty()) {
127
128 if (logger.isInfoEnabled()) {
129 logger.info("No crawling target data configs.");
130 }
131 return;
132 }
133
134 doCrawl(sessionId, configList);
135 }
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153 protected void doCrawl(final String sessionId, final List<DataConfig> configList) {
154 final int multiprocessCrawlingCount = ComponentUtil.getFessConfig().getCrawlingThreadCount();
155
156 final SystemHelper systemHelper = ComponentUtil.getSystemHelper();
157 final long startTime = systemHelper.getCurrentTimeAsLong();
158
159 final IndexUpdateCallback indexUpdateCallback = ComponentUtil.getComponent(IndexUpdateCallback.class);
160
161 final List<String> sessionIdList = new ArrayList<>();
162 dataCrawlingThreadList.clear();
163 final List<String> dataCrawlingThreadStatusList = new ArrayList<>();
164 for (final DataConfig dataConfig : configList) {
165 final DataStoreParams initParamMap = new DataStoreParams();
166 final String sid = ComponentUtil.getCrawlingConfigHelper().store(sessionId, dataConfig);
167 sessionIdList.add(sid);
168
169 initParamMap.put(Constants.SESSION_ID, sessionId);
170 initParamMap.put(Constants.CRAWLING_INFO_ID, sid);
171
172 final DataCrawlingThread dataCrawlingThread = new DataCrawlingThread(dataConfig, indexUpdateCallback, initParamMap);
173 dataCrawlingThread.setPriority(crawlerPriority);
174 dataCrawlingThread.setName(sid);
175 dataCrawlingThread.setDaemon(true);
176
177 dataCrawlingThreadList.add(dataCrawlingThread);
178 dataCrawlingThreadStatusList.add(Constants.READY);
179
180 }
181
182 int startedCrawlerNum = 0;
183 int activeCrawlerNum = 0;
184 while (startedCrawlerNum < dataCrawlingThreadList.size()) {
185
186 if (systemHelper.isForceStop()) {
187 for (final DataCrawlingThread crawlerThread : dataCrawlingThreadList) {
188 crawlerThread.stopCrawling();
189 }
190 break;
191 }
192
193 if (activeCrawlerNum < multiprocessCrawlingCount) {
194
195 dataCrawlingThreadList.get(startedCrawlerNum).start();
196 dataCrawlingThreadStatusList.set(startedCrawlerNum, Constants.RUNNING);
197 startedCrawlerNum++;
198 activeCrawlerNum++;
199 ThreadUtil.sleep(crawlingExecutionInterval);
200 continue;
201 }
202
203
204 for (int i = 0; i < startedCrawlerNum; i++) {
205 if (!dataCrawlingThreadList.get(i).isRunning() && Constants.RUNNING.equals(dataCrawlingThreadStatusList.get(i))) {
206 dataCrawlingThreadList.get(i).awaitTermination();
207 dataCrawlingThreadStatusList.set(i, Constants.DONE);
208 activeCrawlerNum--;
209 }
210 }
211 ThreadUtil.sleep(crawlingExecutionInterval);
212 }
213
214 boolean finishedAll = false;
215 while (!finishedAll) {
216 finishedAll = true;
217 for (int i = 0; i < dataCrawlingThreadList.size(); i++) {
218 dataCrawlingThreadList.get(i).awaitTermination(crawlingExecutionInterval);
219 if (!dataCrawlingThreadList.get(i).isRunning() && Constants.RUNNING.equals(dataCrawlingThreadStatusList.get(i))) {
220 dataCrawlingThreadStatusList.set(i, Constants.DONE);
221 }
222 if (!Constants.DONE.equals(dataCrawlingThreadStatusList.get(i))) {
223 finishedAll = false;
224 }
225 }
226 }
227 dataCrawlingThreadList.clear();
228 dataCrawlingThreadStatusList.clear();
229
230
231 final CrawlingInfoHelper crawlingInfoHelper = ComponentUtil.getCrawlingInfoHelper();
232
233 final long execTime = systemHelper.getCurrentTimeAsLong() - startTime;
234 crawlingInfoHelper.putToInfoMap(Constants.DATA_CRAWLING_EXEC_TIME, Long.toString(execTime));
235 if (logger.isInfoEnabled()) {
236 logger.info("[EXEC TIME] crawling time: {}ms", execTime);
237 }
238
239 crawlingInfoHelper.putToInfoMap(Constants.DATA_INDEX_EXEC_TIME, Long.toString(indexUpdateCallback.getExecuteTime()));
240 crawlingInfoHelper.putToInfoMap(Constants.DATA_INDEX_SIZE, Long.toString(indexUpdateCallback.getDocumentSize()));
241
242 for (final String sid : sessionIdList) {
243
244 ComponentUtil.getCrawlingConfigHelper().remove(sid);
245 }
246
247 }
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262 protected static class DataCrawlingThread extends Thread {
263
264
265 private final DataConfig dataConfig;
266
267
268 private final IndexUpdateCallback indexUpdateCallback;
269
270
271 private final DataStoreParams initParamMap;
272
273
274 protected boolean finished = false;
275
276
277 protected boolean running = false;
278
279
280 private DataStore dataStore;
281
282
283
284
285
286
287
288
289 protected DataCrawlingThread(final DataConfig dataConfig, final IndexUpdateCallback indexUpdateCallback,
290 final DataStoreParams initParamMap) {
291 this.dataConfig = dataConfig;
292 this.indexUpdateCallback = indexUpdateCallback;
293 this.initParamMap = initParamMap;
294 }
295
296
297
298
299
300
301 @Override
302 public void run() {
303 running = true;
304 try {
305 process();
306 } finally {
307 running = false;
308 finished = true;
309 }
310 }
311
312
313
314
315
316
317
318
319 protected void process() {
320 final DataStoreFactory dataStoreFactory = ComponentUtil.getDataStoreFactory();
321 dataStore = dataStoreFactory.getDataStore(dataConfig.getHandlerName());
322 if (dataStore == null) {
323 logger.error("DataStore({}) is not found.", dataConfig.getHandlerName());
324 } else {
325 try {
326 dataStore.store(dataConfig, indexUpdateCallback, initParamMap);
327 } catch (final Throwable e) {
328 logger.error("Failed to process a data crawling: {}", dataConfig.getName(), e);
329 ComponentUtil.getComponent(FailureUrlService.class)
330 .store(dataConfig, e.getClass().getCanonicalName(), dataConfig.getConfigId() + ":" + dataConfig.getName(), e);
331 } finally {
332 indexUpdateCallback.commit();
333 deleteOldDocs();
334 }
335 }
336 }
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352 private void deleteOldDocs() {
353 if (Constants.FALSE.equals(initParamMap.getAsString(DELETE_OLD_DOCS))) {
354 return;
355 }
356 final String sessionId = initParamMap.getAsString(Constants.SESSION_ID);
357 if (StringUtil.isBlank(sessionId)) {
358 logger.warn("[{}] Cannot delete stale documents: sessionId is not set.", dataConfig.getName());
359 return;
360 }
361 final FessConfig fessConfig = ComponentUtil.getFessConfig();
362 final BoolQueryBuilder queryBuilder = QueryBuilders.boolQuery()
363 .must(QueryBuilders.termQuery(fessConfig.getIndexFieldConfigId(), dataConfig.getConfigId()))
364 .mustNot(QueryBuilders.termQuery(fessConfig.getIndexFieldSegment(), sessionId));
365 if (!Constants.FALSE.equals(initParamMap.getAsString(KEEP_EXPIRES_DOCS))) {
366 final QueryBuilder expiresCheckQuery = QueryBuilders.boolQuery()
367 .mustNot(QueryBuilders.rangeQuery(fessConfig.getIndexFieldExpires()).gt("now"))
368 .mustNot(QueryBuilders.existsQuery(fessConfig.getIndexFieldExpires()));
369 queryBuilder.must(expiresCheckQuery);
370 }
371
372 try {
373 final SearchEngineClient searchEngineClient = ComponentUtil.getSearchEngineClient();
374 final String index = fessConfig.getIndexDocumentUpdateIndex();
375 searchEngineClient.admin().indices().prepareRefresh(index).execute().actionGet();
376 final long numOfDeleted = searchEngineClient.deleteByQuery(index, queryBuilder);
377 logger.info("[{}] Deleted {} stale documents.", dataConfig.getName(), numOfDeleted);
378 } catch (final Exception e) {
379 logger.error("[{}] Failed to delete stale documents.", dataConfig.getName(), e);
380 }
381 }
382
383
384
385
386
387
388 public boolean isFinished() {
389 return finished;
390 }
391
392
393
394
395
396
397 public void stopCrawling() {
398 if (dataStore != null) {
399 dataStore.stop();
400 }
401 }
402
403
404
405
406
407
408 public String getCrawlingInfoId() {
409 return initParamMap.getAsString(Constants.CRAWLING_INFO_ID);
410 }
411
412
413
414
415
416
417 public boolean isRunning() {
418 return running;
419 }
420
421
422
423
424
425
426 public void awaitTermination() {
427 try {
428 join();
429 } catch (final InterruptedException e) {
430 if (logger.isDebugEnabled()) {
431 logger.debug("Interrupted.", e);
432 }
433 }
434 }
435
436
437
438
439
440
441
442 public void awaitTermination(final long mills) {
443 try {
444 join(mills);
445 } catch (final InterruptedException e) {
446 if (logger.isDebugEnabled()) {
447 logger.debug("Interrupted.", e);
448 }
449 }
450 }
451 }
452
453
454
455
456
457
458
459
460 public void setCrawlingExecutionInterval(final long crawlingExecutionInterval) {
461 this.crawlingExecutionInterval = crawlingExecutionInterval;
462 }
463
464
465
466
467
468
469
470 public void setCrawlerPriority(final int crawlerPriority) {
471 this.crawlerPriority = crawlerPriority;
472 }
473 }