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.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
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
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
114 if (systemHelper.isForceStop()) {
115 for (final DataCrawlingThread crawlerThread : dataCrawlingThreadList) {
116 crawlerThread.stopCrawling();
117 }
118 break;
119 }
120
121 if (activeCrawlerNum < multiprocessCrawlingCount) {
122
123 dataCrawlingThreadList.get(startedCrawlerNum).start();
124 dataCrawlingThreadStatusList.set(startedCrawlerNum, Constants.RUNNING);
125 startedCrawlerNum++;
126 activeCrawlerNum++;
127 ThreadUtil.sleep(crawlingExecutionInterval);
128 continue;
129 }
130
131
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
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
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 }