Tika Pipes from Java
Tika Pipes is what gives tika-app, tika-server and tika-grpc their
robustness: every parse runs in a forked JVM that the parent feeds over a socket, times out, and
restarts when it dies. The same machinery is available to an application that embeds Tika. This
page walks up the ladder from the simplest entry point to the full pipeline.
| Class | Module | Use when |
|---|---|---|
|
|
You have a file or a stream and want the text and metadata back in your own JVM. Local files and in-memory bytes only. |
|
|
You want one document at a time through a full pipes configuration: any fetcher (S3, HTTP,
a database …), any emitter, results either written by the fork or passed back to you. This is
what |
|
|
You want the whole pipeline: an iterator enumerates documents, a queue feeds a pool of forks,
emitter threads write results, a reporter records status. This is what |
PipesClient is the per-fork socket client both of the last two wrap. It is not a user-facing
API; use the classes above.
Where plugins come from
Fetchers, emitters, iterators and reporters are PF4J plugins. An embedding application can supply them two ways:
-
As Maven dependencies. A plugin jar on the application’s classpath is discovered through its
META-INF/extensions.idx, and the forked JVM inherits the parent’s classpath.tika-pipes-fork-parseralready depends ontika-pipes-file-system, soPipesForkParserworks with no further setup. For another source or destination, add that plugin’s jar (tika-pipes-s3,tika-pipes-opensearch, …) as an ordinary dependency. -
As zips in a plugins directory. This is the distribution layout:
tika-appandtika-servership aplugins/directory of plugin zips, andplugin-rootspoints at it. Get the zips from the distribution zips on the download page; they are not on Maven Central. A zip plugin takes precedence over a classpath plugin of the same name.
Either way, plugin-roots must be present in the JSON config for PipesParser and
AsyncProcessor, even if the directory is empty. PipesForkParser fills it in for you (a
plugins directory beside the jar, else one in the working directory) and logs a warning when it
finds none; with the plugins on the classpath the warning is harmless.
PipesForkParserConfig.setPluginsDir(Path) sets it explicitly.
Dependencies
<!-- PipesForkParser; pulls in tika-pipes-core, the file-system plugin and the standard parsers -->
<dependency>
<groupId>org.apache.tika</groupId>
<artifactId>tika-pipes-fork-parser</artifactId>
<version>4.2.0-SNAPSHOT</version>
</dependency>
<!-- PipesParser and AsyncProcessor without PipesForkParser -->
<dependency>
<groupId>org.apache.tika</groupId>
<artifactId>tika-pipes-core</artifactId>
<version>4.2.0-SNAPSHOT</version>
</dependency>
<dependency>
<groupId>org.apache.tika</groupId>
<artifactId>tika-parsers-standard-package</artifactId>
<version>4.2.0-SNAPSHOT</version>
<type>pom</type>
</dependency>
The forked JVM is launched with the parent’s classpath, so the parsers and plugins your application depends on are the ones the fork uses.
PipesForkParser: one document, result returned
Create the parser once and reuse it; each instance owns a pool of forked JVMs, and creating one per document means starting a JVM per document. It is thread-safe.
import java.nio.file.Path;
import org.apache.tika.config.TimeoutLimits;
import org.apache.tika.metadata.Metadata;
import org.apache.tika.metadata.TikaCoreProperties;
import org.apache.tika.pipes.fork.PipesForkParser;
import org.apache.tika.pipes.fork.PipesForkParserConfig;
import org.apache.tika.pipes.fork.PipesForkResult;
PipesForkParserConfig config = new PipesForkParserConfig()
.setPluginsDir(Path.of("/opt/tika/plugins"))
.setNumClients(2)
.setTimeoutLimits(new TimeoutLimits(
TimeoutLimits.DEFAULT_TOTAL_TASK_TIMEOUT_MILLIS, 60_000))
.addJvmArg("-Xmx1g");
try (PipesForkParser parser = new PipesForkParser(config)) {
for (Path file : files) {
PipesForkResult result = parser.parse(file);
if (result.isSuccess()) {
// container first, then one Metadata per embedded document
for (Metadata m : result.getMetadataList()) {
String text = m.get(TikaCoreProperties.TIKA_CONTENT);
}
} else if (result.isProcessCrash()) {
// OOM or timeout; the fork restarts before the next parse()
} else {
// fetch or parse problem with this document; see result.getMessage()
}
}
}
parse(TikaInputStream) accepts a stream, which is spooled to a temporary file for the fork.
Prefer parse(Path) when you already have a file. result.getContent() is the container’s text
only; iterate getMetadataList() for embedded documents.
See Getting Started with the Java API and PipesForkParserExample.java for parse modes, handler types, content-type hints and error handling.
PipesParser: one document, any source, any destination
PipesParser takes a FetchEmitTuple — a fetch key naming a configured fetcher and a document
within it, plus an emit key naming a configured emitter — and returns a PipesResult. Everything
about the pipeline comes from a JSON config file, so the same code reads from a directory today and
from S3 tomorrow. See Pipeline Configuration for the pipes
block and Plugins for each fetcher and emitter.
{
"plugin-roots": "/opt/tika/plugins",
"fetchers": {
"fsf": {
"file-system-fetcher": { "basePath": "/data/input" }
}
},
"emitters": {
"fse": {
"file-system-emitter": { "basePath": "/data/output", "fileExtension": "json" }
}
},
"pipes": {
"numClients": 2,
"forkedJvmArgs": ["-Xmx1g"],
"emitStrategy": { "type": "EMIT_ALL" }
}
}
import java.nio.file.Path;
import org.apache.tika.pipes.api.FetchEmitTuple;
import org.apache.tika.pipes.api.PipesResult;
import org.apache.tika.pipes.api.emitter.EmitKey;
import org.apache.tika.pipes.api.fetcher.FetchKey;
import org.apache.tika.pipes.core.PipesParser;
try (PipesParser parser = PipesParser.load(Path.of("tika-config.json"))) {
FetchEmitTuple tuple = new FetchEmitTuple(
"doc-1", // your id, echoed in logs and reports
new FetchKey("fsf", "reports/q3.pdf"), // fetcher id, key relative to basePath
new EmitKey("fse", "reports/q3")); // emitter id, key; ".json" is appended
PipesResult result = parser.parse(tuple);
if (result.isSuccess()) {
// EMIT_ALL: the fork wrote /data/output/reports/q3.json; nothing comes back
}
}
Getting the result back instead
To receive the parsed metadata in your JVM, emit nothing and ask for passback:
import org.apache.tika.metadata.Metadata;
import org.apache.tika.parser.ParseContext;
import org.apache.tika.pipes.core.EmitStrategy;
import org.apache.tika.pipes.core.EmitStrategyConfig;
ParseContext context = new ParseContext();
context.set(EmitStrategyConfig.class, new EmitStrategyConfig(EmitStrategy.PASSBACK_ALL));
FetchEmitTuple tuple = new FetchEmitTuple("doc-1",
new FetchKey("fsf", "reports/q3.pdf"), EmitKey.NO_EMIT, new Metadata(), context);
PipesResult result = parser.parse(tuple);
if (result.isSuccess() && result.emitData() != null) {
for (Metadata m : result.emitData().getMetadataList()) {
// container first, then embedded documents
}
}
The default emitStrategy is DYNAMIC: small extracts are passed back, large ones are written
by the fork. With PipesParser there is no emitter thread behind you, so a passed-back result
is yours to write; set EMIT_ALL when the fork should write everything, PASSBACK_ALL when you
want everything back. The per-request ParseContext setting above overrides the config.
PipesParser is thread-safe. Each call borrows one of numClients forks and blocks for up to
maxWaitForClientMillis for one to free up, returning CLIENT_UNAVAILABLE_WITHIN_MS if none
does.
AsyncProcessor: the whole pipeline
AsyncProcessor adds the queue, the worker threads, the emitter threads and the reporter. You
offer FetchEmitTuple`s and it does the rest. Add a `pipes-iterator and a reporter to the config
above and the tuples come from the config too:
"pipes-iterator": {
"file-system-pipes-iterator": {
"basePath": "/data/input",
"fetcherId": "fsf",
"emitterId": "fse"
}
},
"pipes-reporters": {
"file-system-reporter": { "statusFile": "/data/status.json" }
}
import java.nio.file.Path;
import java.util.concurrent.TimeoutException;
import org.apache.tika.config.loader.TikaJsonConfig;
import org.apache.tika.pipes.api.FetchEmitTuple;
import org.apache.tika.pipes.api.pipesiterator.PipesIterator;
import org.apache.tika.pipes.core.async.AsyncProcessor;
import org.apache.tika.pipes.core.pipesiterator.PipesIteratorManager;
import org.apache.tika.plugins.TikaPluginManager;
Path configPath = Path.of("tika-config.json");
TikaJsonConfig jsonConfig = TikaJsonConfig.load(configPath);
PipesIterator iterator = PipesIteratorManager
.load(TikaPluginManager.load(jsonConfig), jsonConfig)
.orElseThrow();
try (AsyncProcessor processor = AsyncProcessor.load(configPath, iterator)) {
for (FetchEmitTuple tuple : iterator) {
if (!processor.offer(tuple, 120_000)) {
throw new TimeoutException("queue full");
}
}
processor.finished(); // no more input
while (processor.checkActive()) { // rethrows the first worker or emitter failure
Thread.sleep(500);
}
long processed = processor.getTotalProcessed();
}
You do not need an iterator plugin. Pass null as the second argument to load and offer tuples
you build yourself, from a message queue, a database cursor or anything else. offer blocks up
to the given milliseconds while the queue is full (queueSize in the pipes block), and returns
false if it never drained. Per-document outcomes go to the configured
reporter, not back to the caller; checkActive only surfaces failures
that stop the pipeline. AsyncProcessor also handles emission of passed-back extracts, so the
default DYNAMIC emit strategy is the right one here.
Reading a result
PipesResult (and PipesForkResult, which wraps one) sorts every status into a category, with a
predicate for each:
| Predicate | Statuses | What to do |
|---|---|---|
|
|
Use the output. A |
|
|
This document failed; the fork is fine. Log it and move on. |
|
|
The fork died on this document and restarts before the next call. Record the document; do not retry it on the same settings. |
|
|
Possibly transient: a backend is down, or every fork was busy. Retry later. |
|
|
The fork cannot start at all: bad config, bad classpath, bad |
PipesForkParser throws PipesForkParserException for configuration and initialization
problems instead of returning them.
Lifecycle
All three classes are Closeable. Closing kills the forked JVMs and removes the temporary files
the parent created for them. Each fork also restarts itself after maxFilesProcessedPerProcess documents to
bound slow leaks in the parsing libraries. Sizing the pool — numClients, heap per fork, CPU
per fork — is covered in Forked-JVM CPU and Heap Sizing, and the
timeout model in Timeouts.