Author: thomasm
Date: Thu Nov 9 16:52:49 2017
New Revision: 1814745
URL: http://svn.apache.org/viewvc?rev=1814745&view=rev
Log:
OAK-5519 Skip problematic binaries instead of blocking indexing
Modified:
jackrabbit/oak/trunk/oak-lucene/src/main/java/org/apache/jackrabbit/oak/plugins/index/lucene/ExtractedTextCache.java
jackrabbit/oak/trunk/oak-lucene/src/main/java/org/apache/jackrabbit/oak/plugins/index/lucene/LuceneIndexProviderService.java
jackrabbit/oak/trunk/oak-lucene/src/main/java/org/apache/jackrabbit/oak/plugins/index/lucene/TextExtractionStatsMBean.java
jackrabbit/oak/trunk/oak-lucene/src/main/java/org/apache/jackrabbit/oak/plugins/index/lucene/binary/BinaryTextExtractor.java
jackrabbit/oak/trunk/oak-lucene/src/test/java/org/apache/jackrabbit/oak/plugins/index/lucene/ExtractedTextCacheTest.java
Modified:
jackrabbit/oak/trunk/oak-lucene/src/main/java/org/apache/jackrabbit/oak/plugins/index/lucene/ExtractedTextCache.java
URL:
http://svn.apache.org/viewvc/jackrabbit/oak/trunk/oak-lucene/src/main/java/org/apache/jackrabbit/oak/plugins/index/lucene/ExtractedTextCache.java?rev=1814745&r1=1814744&r2=1814745&view=diff
==============================================================================
---
jackrabbit/oak/trunk/oak-lucene/src/main/java/org/apache/jackrabbit/oak/plugins/index/lucene/ExtractedTextCache.java
(original)
+++
jackrabbit/oak/trunk/oak-lucene/src/main/java/org/apache/jackrabbit/oak/plugins/index/lucene/ExtractedTextCache.java
Thu Nov 9 16:52:49 2017
@@ -19,8 +19,23 @@
package org.apache.jackrabbit.oak.plugins.index.lucene;
+import java.io.File;
+import java.io.FileInputStream;
+import java.io.FileOutputStream;
import java.io.IOException;
+import java.util.Map.Entry;
+import java.util.Properties;
+import java.util.concurrent.Callable;
+import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.ExecutionException;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Future;
+import java.util.concurrent.LinkedBlockingQueue;
+import java.util.concurrent.ThreadFactory;
+import java.util.concurrent.ThreadPoolExecutor;
import java.util.concurrent.TimeUnit;
+import java.util.concurrent.TimeoutException;
+import java.util.concurrent.atomic.AtomicInteger;
import javax.annotation.CheckForNull;
import javax.annotation.Nonnull;
@@ -32,6 +47,7 @@ import org.apache.jackrabbit.oak.api.Blo
import org.apache.jackrabbit.oak.cache.CacheStats;
import org.apache.jackrabbit.oak.commons.IOUtils;
import org.apache.jackrabbit.oak.plugins.index.fulltext.ExtractedText;
+import
org.apache.jackrabbit.oak.plugins.index.fulltext.ExtractedText.ExtractionResult;
import
org.apache.jackrabbit.oak.plugins.index.fulltext.PreExtractedTextProvider;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -39,6 +55,18 @@ import org.slf4j.LoggerFactory;
import static org.apache.jackrabbit.oak.commons.PathUtils.concat;
public class ExtractedTextCache {
+ private static final boolean CACHE_ONLY_SUCCESS =
+ Boolean.getBoolean("oak.extracted.cacheOnlySuccess");
+ private static final int EXTRACTION_TIMEOUT_SECONDS =
+ Integer.getInteger("oak.extraction.timeoutSeconds", 60);
+ private static final int EXTRACTION_MAX_THREADS =
+ Integer.getInteger("oak.extraction.maxThreads", 10);
+ private static final boolean EXTRACT_IN_CALLER_THREAD =
+ Boolean.getBoolean("oak.extraction.inCallerThread");
+ private static final boolean EXTRACT_FORGET_TIMEOUT =
+ Boolean.getBoolean("oak.extraction.forgetTimeout");
+
+ private static final String TIMEOUT_MAP =
"textExtractionTimeout.properties";
private static final String EMPTY_STRING = "";
private static final Logger log =
LoggerFactory.getLogger(ExtractedTextCache.class);
private volatile PreExtractedTextProvider extractedTextProvider;
@@ -48,14 +76,20 @@ public class ExtractedTextCache {
private long totalTime;
private int preFetchedCount;
private final Cache<String, String> cache;
+ private final ConcurrentHashMap<String, String> timeoutMap;
+ private final File indexDir;
private final CacheStats cacheStats;
private final boolean alwaysUsePreExtractedCache;
+ private volatile ExecutorService executorService;
+ private volatile int timeoutCount;
+ private long extractionTimeoutMillis = EXTRACTION_TIMEOUT_SECONDS * 1000;
public ExtractedTextCache(long maxWeight, long expiryTimeInSecs){
- this(maxWeight, expiryTimeInSecs, false);
+ this(maxWeight, expiryTimeInSecs, false, null);
}
- public ExtractedTextCache(long maxWeight, long expiryTimeInSecs, boolean
alwaysUsePreExtractedCache) {
+ public ExtractedTextCache(long maxWeight, long expiryTimeInSecs, boolean
alwaysUsePreExtractedCache,
+ File indexDir) {
if (maxWeight > 0) {
cache = CacheBuilder.newBuilder()
.weigher(EmpiricalWeigher.INSTANCE)
@@ -70,6 +104,9 @@ public class ExtractedTextCache {
cacheStats = null;
}
this.alwaysUsePreExtractedCache = alwaysUsePreExtractedCache;
+ this.timeoutMap = new ConcurrentHashMap<String, String>();
+ this.indexDir = indexDir;
+ loadTimeoutMap();
}
/**
@@ -90,39 +127,53 @@ public class ExtractedTextCache {
ExtractedText text =
extractedTextProvider.getText(propertyPath, blob);
if (text != null) {
preFetchedCount++;
- switch (text.getExtractionResult()) {
- case SUCCESS:
- result = text.getExtractedText().toString();
- break;
- case ERROR:
- result = LuceneIndexEditor.TEXT_EXTRACTION_ERROR;
- break;
- case EMPTY:
- result = EMPTY_STRING;
- break;
- }
+ result = getText(text);
}
} catch (IOException e) {
log.warn("Error occurred while fetching pre extracted text for
{}", propertyPath, e);
}
}
-
String id = blob.getContentIdentity();
if (cache != null && id != null && result == null) {
result = cache.getIfPresent(id);
}
+ if (result == null && id != null) {
+ result = timeoutMap.get(id);
+ }
return result;
}
public void put(@Nonnull Blob blob, @Nonnull ExtractedText extractedText) {
String id = blob.getContentIdentity();
- if (extractedText.getExtractionResult() ==
ExtractedText.ExtractionResult.SUCCESS
- && cache != null
- && id != null) {
- cache.put(id, extractedText.getExtractedText().toString());
+ if (cache != null && id != null) {
+ if (extractedText.getExtractionResult() == ExtractionResult.SUCCESS
+ || !CACHE_ONLY_SUCCESS) {
+ cache.put(id, getText(extractedText));
+ }
}
}
+ public void putTimeout(@Nonnull Blob blob, @Nonnull ExtractedText
extractedText) {
+ if (EXTRACT_FORGET_TIMEOUT) {
+ return;
+ }
+ String id = blob.getContentIdentity();
+ timeoutMap.put(id, getText(extractedText));
+ storeTimeoutMap();
+ }
+
+ private static String getText(ExtractedText text) {
+ switch (text.getExtractionResult()) {
+ case SUCCESS:
+ return text.getExtractedText().toString();
+ case ERROR:
+ return LuceneIndexEditor.TEXT_EXTRACTION_ERROR;
+ case EMPTY:
+ return EMPTY_STRING;
+ }
+ throw new IllegalArgumentException();
+ }
+
public void addStats(int count, long timeInMillis, long bytesRead, long
textLength){
this.textExtractionCount += count;
this.totalTime += timeInMillis;
@@ -130,7 +181,7 @@ public class ExtractedTextCache {
this.totalTextSize += textLength;
}
- public TextExtractionStatsMBean getStatsMBean(){
+ public TextExtractionStatsMBean getStatsMBean() {
return new TextExtractionStatsMBean() {
@Override
public boolean isPreExtractedTextProviderConfigured() {
@@ -166,6 +217,11 @@ public class ExtractedTextCache {
public boolean isAlwaysUsePreExtractedCache() {
return alwaysUsePreExtractedCache;
}
+
+ @Override
+ public int getTimeoutCount() {
+ return timeoutCount;
+ }
};
}
@@ -216,4 +272,133 @@ public class ExtractedTextCache {
return (int) size;
}
}
+
+ public void close() {
+ resetCache();
+ // don't clean the persistent map on purpose, so we don't re-try
+ // after restarting the service or so
+ closeExecutorService();
+ }
+
+ public void process(String name, Callable<Void> callable) throws
InterruptedException, Throwable {
+ Callable<Void> callable2 = new Callable<Void>() {
+ @Override
+ public Void call() throws Exception {
+ Thread t = Thread.currentThread();
+ String oldThreadName = t.getName();
+ t.setName(oldThreadName + ": " + name);
+ try {
+ return callable.call();
+ } finally {
+ Thread.currentThread().setName(oldThreadName);
+ }
+ }
+ };
+ try {
+ if (EXTRACT_IN_CALLER_THREAD) {
+ callable2.call();
+ } else {
+ Future<Void> future = getExecutor().submit(callable2);
+ future.get(extractionTimeoutMillis, TimeUnit.MILLISECONDS);
+ }
+ } catch (TimeoutException e) {
+ timeoutCount++;
+ throw e;
+ } catch (InterruptedException e) {
+ throw e;
+ } catch (ExecutionException e) {
+ throw e.getCause();
+ }
+ }
+
+ public void setExtractionTimeoutMillis(int extractionTimeoutMillis) {
+ this.extractionTimeoutMillis = extractionTimeoutMillis;
+ }
+
+ private ExecutorService getExecutor() {
+ if (executorService == null) {
+ createExecutor();
+ }
+ return executorService;
+ }
+
+ private synchronized void createExecutor() {
+ if (executorService != null) {
+ return;
+ }
+ log.debug("ExtractedTextCache createExecutor " + this);
+ ThreadPoolExecutor executor = new ThreadPoolExecutor(1,
EXTRACTION_MAX_THREADS,
+ 60L, TimeUnit.SECONDS,
+ new LinkedBlockingQueue<Runnable>(), new ThreadFactory() {
+ private final AtomicInteger counter = new AtomicInteger();
+ private final Thread.UncaughtExceptionHandler handler = new
Thread.UncaughtExceptionHandler() {
+ @Override
+ public void uncaughtException(Thread t, Throwable e) {
+ log.warn("Error occurred in asynchronous processing ", e);
+ }
+ };
+ @Override
+ public Thread newThread(@Nonnull Runnable r) {
+ Thread thread = new Thread(r, createName());
+ thread.setDaemon(true);
+ thread.setPriority(Thread.MIN_PRIORITY);
+ thread.setUncaughtExceptionHandler(handler);
+ return thread;
+ }
+
+ private String createName() {
+ int index = counter.getAndIncrement();
+ return "oak binary text extractor" + (index == 0 ? "" : " " +
index);
+ }
+ });
+ executor.setKeepAliveTime(1, TimeUnit.MINUTES);
+ executor.allowCoreThreadTimeOut(true);
+ executorService = executor;
+ }
+
+ private synchronized void closeExecutorService() {
+ if (executorService != null) {
+ log.debug("ExtractedTextCache closeExecutorService " + this);
+ executorService.shutdown();
+ try {
+ executorService.awaitTermination(1, TimeUnit.MINUTES);
+ } catch (InterruptedException e) {
+ log.warn("Interrupted", e);
+ }
+ executorService = null;
+ }
+ }
+
+ private synchronized void loadTimeoutMap() {
+ if (indexDir == null || !indexDir.exists()) {
+ return;
+ }
+ try (FileInputStream in = new FileInputStream(
+ new File(indexDir, TIMEOUT_MAP))) {
+ Properties prop = new Properties();
+ prop.load(in);
+ for(Entry<Object, Object> e : prop.entrySet()) {
+ timeoutMap.put(e.getKey().toString(), e.getValue().toString());
+ }
+ } catch (Exception e) {
+ log.warn("Could not load timeout map {} from {}",
+ TIMEOUT_MAP, indexDir, e);
+ }
+ }
+
+ private synchronized void storeTimeoutMap() {
+ if (indexDir == null || !indexDir.exists()) {
+ return;
+ }
+ try (FileOutputStream out = new FileOutputStream(
+ new File(indexDir, TIMEOUT_MAP))) {
+ Properties prop = new Properties();
+ prop.putAll(timeoutMap);
+ prop.store(out, "Text extraction timed out for the following
binaries, and will not be retried");
+ } catch (Exception e) {
+ log.warn("Could not store timeout map {} from {}",
+ TIMEOUT_MAP, indexDir, e);
+ }
+ }
+
}
Modified:
jackrabbit/oak/trunk/oak-lucene/src/main/java/org/apache/jackrabbit/oak/plugins/index/lucene/LuceneIndexProviderService.java
URL:
http://svn.apache.org/viewvc/jackrabbit/oak/trunk/oak-lucene/src/main/java/org/apache/jackrabbit/oak/plugins/index/lucene/LuceneIndexProviderService.java?rev=1814745&r1=1814744&r2=1814745&view=diff
==============================================================================
---
jackrabbit/oak/trunk/oak-lucene/src/main/java/org/apache/jackrabbit/oak/plugins/index/lucene/LuceneIndexProviderService.java
(original)
+++
jackrabbit/oak/trunk/oak-lucene/src/main/java/org/apache/jackrabbit/oak/plugins/index/lucene/LuceneIndexProviderService.java
Thu Nov 9 16:52:49 2017
@@ -434,6 +434,10 @@ public class LuceneIndexProviderService
executorService.awaitTermination(1, TimeUnit.MINUTES);
}
+ if (extractedTextCache != null) {
+ extractedTextCache.close();
+ }
+
InfoStream.setDefault(InfoStream.NO_OUTPUT);
}
@@ -673,7 +677,11 @@ public class LuceneIndexProviderService
boolean alwaysUsePreExtractedCache =
PropertiesUtil.toBoolean(config.get(PROP_PRE_EXTRACTED_TEXT_ALWAYS_USE),
PROP_PRE_EXTRACTED_TEXT_ALWAYS_USE_DEFAULT);
- extractedTextCache = new ExtractedTextCache(cacheSizeInMB * ONE_MB,
cacheExpiryInSecs, alwaysUsePreExtractedCache);
+ extractedTextCache = new ExtractedTextCache(
+ cacheSizeInMB * ONE_MB,
+ cacheExpiryInSecs,
+ alwaysUsePreExtractedCache,
+ indexDir);
if (extractedTextProvider != null){
registerExtractedTextProvider(extractedTextProvider);
}
Modified:
jackrabbit/oak/trunk/oak-lucene/src/main/java/org/apache/jackrabbit/oak/plugins/index/lucene/TextExtractionStatsMBean.java
URL:
http://svn.apache.org/viewvc/jackrabbit/oak/trunk/oak-lucene/src/main/java/org/apache/jackrabbit/oak/plugins/index/lucene/TextExtractionStatsMBean.java?rev=1814745&r1=1814744&r2=1814745&view=diff
==============================================================================
---
jackrabbit/oak/trunk/oak-lucene/src/main/java/org/apache/jackrabbit/oak/plugins/index/lucene/TextExtractionStatsMBean.java
(original)
+++
jackrabbit/oak/trunk/oak-lucene/src/main/java/org/apache/jackrabbit/oak/plugins/index/lucene/TextExtractionStatsMBean.java
Thu Nov 9 16:52:49 2017
@@ -38,4 +38,6 @@ public interface TextExtractionStatsMBea
String getExtractedTextSize();
String getBytesRead();
+
+ int getTimeoutCount();
}
Modified:
jackrabbit/oak/trunk/oak-lucene/src/main/java/org/apache/jackrabbit/oak/plugins/index/lucene/binary/BinaryTextExtractor.java
URL:
http://svn.apache.org/viewvc/jackrabbit/oak/trunk/oak-lucene/src/main/java/org/apache/jackrabbit/oak/plugins/index/lucene/binary/BinaryTextExtractor.java?rev=1814745&r1=1814744&r2=1814745&view=diff
==============================================================================
---
jackrabbit/oak/trunk/oak-lucene/src/main/java/org/apache/jackrabbit/oak/plugins/index/lucene/binary/BinaryTextExtractor.java
(original)
+++
jackrabbit/oak/trunk/oak-lucene/src/main/java/org/apache/jackrabbit/oak/plugins/index/lucene/binary/BinaryTextExtractor.java
Thu Nov 9 16:52:49 2017
@@ -26,6 +26,8 @@ import java.util.ArrayList;
import java.util.Collections;
import java.util.List;
import java.util.Set;
+import java.util.concurrent.Callable;
+import java.util.concurrent.TimeoutException;
import javax.annotation.Nullable;
@@ -140,16 +142,21 @@ public class BinaryTextExtractor {
if (log.isDebugEnabled()) {
log.debug("Extracting {}, {} bytes, id {}", path, length,
v.getContentIdentity());
}
- String oldThreadName = null;
- if (length > SMALL_BINARY) {
- Thread t = Thread.currentThread();
- oldThreadName = t.getName();
- t.setName(oldThreadName + ": Extracting " + path + ", " + length +
" bytes");
- }
try {
CountingInputStream stream = new CountingInputStream(new
LazyInputStream(new BlobByteSource(v)));
try {
- getParser().parse(stream, handler, metadata, new
ParseContext());
+ if (length > SMALL_BINARY) {
+ String name = "Extracting " + path + ", " + length + "
bytes";
+ extractedTextCache.process(name, new Callable<Void>() {
+ @Override
+ public Void call() throws Exception {
+ getParser().parse(stream, handler, metadata, new
ParseContext());
+ return null;
+ }
+ });
+ } else {
+ getParser().parse(stream, handler, metadata, new
ParseContext());
+ }
} finally {
bytesRead = stream.getCount();
stream.close();
@@ -159,6 +166,13 @@ public class BinaryTextExtractor {
// not being present. This is equivalent to disabling
// selected media types in configuration, so we can simply
// ignore these errors.
+ } catch (TimeoutException t) {
+ log.warn(
+ "[{}] Failed to extract text from a binary property due to
timeout: {}.",
+ getIndexName(), path);
+ extractedTextCache.put(v, ExtractedText.ERROR);
+ extractedTextCache.putTimeout(v, ExtractedText.ERROR);
+ return TEXT_EXTRACTION_ERROR;
} catch (Throwable t) {
// Capture and report any other full text extraction problems.
// The special STOP exception is used for normal termination.
@@ -172,10 +186,6 @@ public class BinaryTextExtractor {
extractedTextCache.put(v, ExtractedText.ERROR);
return TEXT_EXTRACTION_ERROR;
}
- } finally {
- if (oldThreadName != null) {
- Thread.currentThread().setName(oldThreadName);
- }
}
String result = handler.toString();
if (bytesRead > 0) {
Modified:
jackrabbit/oak/trunk/oak-lucene/src/test/java/org/apache/jackrabbit/oak/plugins/index/lucene/ExtractedTextCacheTest.java
URL:
http://svn.apache.org/viewvc/jackrabbit/oak/trunk/oak-lucene/src/test/java/org/apache/jackrabbit/oak/plugins/index/lucene/ExtractedTextCacheTest.java?rev=1814745&r1=1814744&r2=1814745&view=diff
==============================================================================
---
jackrabbit/oak/trunk/oak-lucene/src/test/java/org/apache/jackrabbit/oak/plugins/index/lucene/ExtractedTextCacheTest.java
(original)
+++
jackrabbit/oak/trunk/oak-lucene/src/test/java/org/apache/jackrabbit/oak/plugins/index/lucene/ExtractedTextCacheTest.java
Thu Nov 9 16:52:49 2017
@@ -27,15 +27,20 @@ import org.apache.jackrabbit.oak.plugins
import org.apache.jackrabbit.oak.plugins.memory.ArrayBasedBlob;
import org.junit.Test;
+import static org.junit.Assert.fail;
import static org.junit.Assert.assertEquals;
import static org.junit.Assert.assertNotNull;
import static org.junit.Assert.assertNull;
+import static org.junit.Assert.assertTrue;
import static org.mockito.Matchers.any;
import static org.mockito.Matchers.anyString;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.verifyZeroInteractions;
import static org.mockito.Mockito.when;
+import java.util.concurrent.Callable;
+import java.util.concurrent.TimeoutException;
+
public class ExtractedTextCacheTest {
@Test
@@ -84,7 +89,7 @@ public class ExtractedTextCacheTest {
cache.put(b, new ExtractedText(ExtractionResult.ERROR, "test hello"));
text = cache.get("/a", "foo", b, false);
- assertNull(text);
+ assertEquals(LuceneIndexEditor.TEXT_EXTRACTION_ERROR, text);
}
@Test
@@ -124,7 +129,7 @@ public class ExtractedTextCacheTest {
@Test
public void preExtractionAlwaysUse() throws Exception{
- ExtractedTextCache cache = new ExtractedTextCache(10 *
FileUtils.ONE_MB, 100, true);
+ ExtractedTextCache cache = new ExtractedTextCache(10 *
FileUtils.ONE_MB, 100, true, null);
PreExtractedTextProvider provider =
mock(PreExtractedTextProvider.class);
cache.setExtractedTextProvider(provider);
@@ -135,6 +140,51 @@ public class ExtractedTextCacheTest {
assertEquals("bar", text);
}
+ @Test
+ public void rememberTimeout() throws Exception{
+ ExtractedTextCache cache = new ExtractedTextCache(0, 0, false, null);
+ Blob b = new IdBlob("hello", "a");
+ cache.put(b, ExtractedText.ERROR);
+ assertNull(cache.get("/a", "foo", b, false));
+ cache.putTimeout(b, ExtractedText.ERROR);
+ assertEquals(LuceneIndexEditor.TEXT_EXTRACTION_ERROR, cache.get("/a",
"foo", b, false));
+ }
+
+ @Test
+ public void process() throws Throwable {
+ ExtractedTextCache cache = new ExtractedTextCache(0, 0, false, null);
+ try {
+ cache.process("test", new Callable<Void>() {
+ @Override
+ public Void call() throws Exception {
+ throw new OutOfMemoryError();
+ }
+ });
+ fail();
+ } catch (OutOfMemoryError e) {
+ // expected
+ }
+ assertEquals(0, cache.getStatsMBean().getTimeoutCount());
+ cache.setExtractionTimeoutMillis(10);
+ long time = System.currentTimeMillis();
+ try {
+ cache.process("test", new Callable<Void>() {
+ @Override
+ public Void call() throws Exception {
+ // this happens in the background, so doesn't block the
test
+ Thread.sleep(10000);
+ return null;
+ }
+ });
+ fail();
+ } catch (TimeoutException e) {
+ // expected
+ }
+ time = System.currentTimeMillis() - time;
+ assertTrue("" + time, time < 5000);
+ assertEquals(1, cache.getStatsMBean().getTimeoutCount());
+ }
+
private static class IdBlob extends ArrayBasedBlob {
final String id;