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

PipesForkParser

tika-pipes-fork-parser

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.

PipesParser

tika-pipes-core

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 tika-server runs behind /tika and /rmeta.

AsyncProcessor

tika-pipes-core

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 tika-app -i/-o runs.

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-parser already depends on tika-pipes-file-system, so PipesForkParser works 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-app and tika-server ship a plugins/ directory of plugin zips, and plugin-roots points 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.1.0</version>
</dependency>

<!-- PipesParser and AsyncProcessor without PipesForkParser -->
<dependency>
    <groupId>org.apache.tika</groupId>
    <artifactId>tika-pipes-core</artifactId>
    <version>4.1.0</version>
</dependency>
<dependency>
    <groupId>org.apache.tika</groupId>
    <artifactId>tika-parsers-standard-package</artifactId>
    <version>4.1.0</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

isSuccess()

PARSE_SUCCESS, EMIT_SUCCESS, EMIT_SUCCESS_PASSBACK, PARSE_SUCCESS_WITH_EXCEPTION, EMIT_SUCCESS_PARSE_EXCEPTION, PARTIAL_TIMEOUT, EMPTY_OUTPUT

Use the output. A *_EXCEPTION variant carries a partial parse, with the stack trace in message() and in TikaCoreProperties.CONTAINER_EXCEPTION; PARTIAL_TIMEOUT means the task deadline cut the parse short.

isTaskException()

FETCH_EXCEPTION, EMIT_EXCEPTION, FETCHER_NOT_FOUND, EMITTER_NOT_FOUND, PRESET_NOT_FOUND, PAYLOAD_LIMIT_EXCEEDED

This document failed; the fork is fine. Log it and move on.

isProcessCrash()

OOM, TIMEOUT, UNSPECIFIED_CRASH

The fork died on this document and restarts before the next call. Record the document; do not retry it on the same settings.

isInitializationFailure()

FETCHER_INITIALIZATION_EXCEPTION, EMITTER_INITIALIZATION_EXCEPTION, CLIENT_UNAVAILABLE_WITHIN_MS

Possibly transient: a backend is down, or every fork was busy. Retry later.

isFatal()

FAILED_TO_INITIALIZE

The fork cannot start at all: bad config, bad classpath, bad javaPath. Stop and fix it.

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.