This is an automated email from the ASF dual-hosted git repository.
tballison pushed a commit to branch docs/4.0.x
in repository https://gitbox.apache.org/repos/asf/tika.git
The following commit(s) were added to refs/heads/docs/4.0.x by this push:
new 988eb0ee6c improve pipes documentation (backport to 4.0.x) (#3224)
988eb0ee6c is described below
commit 988eb0ee6c9a774e9f3e5e5b081f0f992464f261
Author: Tim Allison <[email protected]>
AuthorDate: Wed Sep 23 16:09:02 2026 -0400
improve pipes documentation (backport to 4.0.x) (#3224)
---
docs/modules/ROOT/nav.adoc | 1 +
docs/modules/ROOT/pages/pipes/configuration.adoc | 2 +-
docs/modules/ROOT/pages/pipes/getting-started.adoc | 2 +-
.../pages/using-tika/java-api/getting-started.adoc | 3 +-
.../ROOT/pages/using-tika/java-api/index.adoc | 3 +-
.../ROOT/pages/using-tika/java-api/pipes.adoc | 338 +++++++++++++++++++++
6 files changed, 345 insertions(+), 4 deletions(-)
diff --git a/docs/modules/ROOT/nav.adoc b/docs/modules/ROOT/nav.adoc
index 76e27d587f..2f0a10f9da 100644
--- a/docs/modules/ROOT/nav.adoc
+++ b/docs/modules/ROOT/nav.adoc
@@ -18,6 +18,7 @@
* xref:using-tika/index.adoc[Using Tika]
** xref:using-tika/java-api/index.adoc[Java API]
*** xref:using-tika/java-api/getting-started.adoc[Getting Started with the
Java API]
+*** xref:using-tika/java-api/pipes.adoc[Tika Pipes from Java]
** xref:using-tika/cli/index.adoc[Command Line]
** xref:using-tika/server/index.adoc[Tika Server]
*** xref:using-tika/server/tls.adoc[TLS/SSL Configuration]
diff --git a/docs/modules/ROOT/pages/pipes/configuration.adoc
b/docs/modules/ROOT/pages/pipes/configuration.adoc
index 02a309661e..a18a3e26e3 100644
--- a/docs/modules/ROOT/pages/pipes/configuration.adoc
+++ b/docs/modules/ROOT/pages/pipes/configuration.adoc
@@ -41,7 +41,7 @@ how many forked JVMs to run, timeouts, memory management, and
parse behavior.
|`numClients`
|_CPU-derived_
-|Number of parallel forked JVMs. Each processes one document at a time.
Defaults to `max(1, min((cores - 2) / 2, 4))` -- roughly half the host's cores,
capped at 4. The batch CLIs (`tika-app -i/-o`, `tika-async-cli`) override that
to `2` when neither `-n` nor a config value is given -- see
xref:using-tika/cli/index.adoc#_tika_pipes_processing[Tika Pipes Processing].
See xref:pipes/cpu-sizing.adoc[Forked-JVM CPU and Heap Sizing] for guidance on
choosing this value relative to host CPU count.
+|Number of parallel forked JVMs. Each processes one document at a time.
Defaults to `+max(1, min((cores - 2) / 2, 4))+` -- roughly half the host's
cores, capped at 4. The batch CLIs (`tika-app -i/-o`, `tika-async-cli`)
override that to `2` when neither `-n` nor a config value is given -- see
xref:using-tika/cli/index.adoc#_tika_pipes_processing[Tika Pipes Processing].
See xref:pipes/cpu-sizing.adoc[Forked-JVM CPU and Heap Sizing] for guidance on
choosing this value relative to host CPU count.
|`forkedJvmArgs`
|`[]`
diff --git a/docs/modules/ROOT/pages/pipes/getting-started.adoc
b/docs/modules/ROOT/pages/pipes/getting-started.adoc
index c66f1d2320..eb83c0a801 100644
--- a/docs/modules/ROOT/pages/pipes/getting-started.adoc
+++ b/docs/modules/ROOT/pages/pipes/getting-started.adoc
@@ -110,7 +110,7 @@ The `pipes` section controls the pipeline. The four
settings worth knowing first
|`numClients`
|_CPU-derived_
-|Number of parallel forked parse processes. `max(1, min((cores - 2) / 2, 4))`
— roughly half the host's cores, capped at 4.
+|Number of parallel forked parse processes. `+max(1, min((cores - 2) / 2,
4))+` — roughly half the host's cores, capped at 4.
|`parseMode`
|`RMETA`
diff --git a/docs/modules/ROOT/pages/using-tika/java-api/getting-started.adoc
b/docs/modules/ROOT/pages/using-tika/java-api/getting-started.adoc
index 550c212cba..289f385f0e 100644
--- a/docs/modules/ROOT/pages/using-tika/java-api/getting-started.adoc
+++ b/docs/modules/ROOT/pages/using-tika/java-api/getting-started.adoc
@@ -66,7 +66,8 @@ try (PipesForkParser parser = new PipesForkParser()) {
See
https://github.com/apache/tika/blob/main/tika-example/src/main/java/org/apache/tika/example/PipesForkParserExample.java[PipesForkParserExample.java]
in the `tika-example` module for embedded documents, custom configuration,
error handling and batch
-processing.
+processing. `PipesForkParser` reads local files only; for other sources and
destinations, or for
+a whole iterator-to-emitter pipeline, see
xref:using-tika/java-api/pipes.adoc[Tika Pipes from Java].
== Without process isolation
diff --git a/docs/modules/ROOT/pages/using-tika/java-api/index.adoc
b/docs/modules/ROOT/pages/using-tika/java-api/index.adoc
index 57db4b8bcb..b0072298e6 100644
--- a/docs/modules/ROOT/pages/using-tika/java-api/index.adoc
+++ b/docs/modules/ROOT/pages/using-tika/java-api/index.adoc
@@ -30,7 +30,8 @@ xref:pipes/index.adoc[Tika Pipes], which runs each parse in a
forked JVM with ti
limits; xref:using-tika/server/index.adoc[tika-server] and
xref:using-tika/grpc/index.adoc[tika-grpc]
provide the same robustness as a service. See
xref:using-tika/java-api/getting-started.adoc[Getting Started with the Java
API] for how to choose,
-and xref:advanced/robustness.adoc[Robustness] for what goes wrong without
isolation.
+xref:using-tika/java-api/pipes.adoc[Tika Pipes from Java] for the Pipes API,
and
+xref:advanced/robustness.adoc[Robustness] for what goes wrong without
isolation.
Know where the line falls before you call `AutoDetectParser` in your own JVM.
Tika's sandboxing is
the forked JVM, and calling a parser directly opts out of it. Untrusted data
can often
diff --git a/docs/modules/ROOT/pages/using-tika/java-api/pipes.adoc
b/docs/modules/ROOT/pages/using-tika/java-api/pipes.adoc
new file mode 100644
index 0000000000..27eaecf9bf
--- /dev/null
+++ b/docs/modules/ROOT/pages/using-tika/java-api/pipes.adoc
@@ -0,0 +1,338 @@
+//
+// Licensed to the Apache Software Foundation (ASF) under one or more
+// contributor license agreements. See the NOTICE file distributed with
+// this work for additional information regarding copyright ownership.
+// The ASF licenses this file to You under the Apache License, Version 2.0
+// (the "License"); you may not use this file except in compliance with
+// the License. You may obtain a copy of the License at
+//
+// http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing, software
+// distributed under the License is distributed on an "AS IS" BASIS,
+// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+// See the License for the specific language governing permissions and
+// limitations under the License.
+//
+
+= Tika Pipes from Java
+:toc:
+:toclevels: 2
+
+xref:pipes/index.adoc[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.
+
+[cols="1,1,3"]
+|===
+|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
xref:pipes/plugins/index.adoc[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
https://tika.apache.org/download.html[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
+
+[source,xml,subs=attributes+]
+----
+<!-- 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>{tika-version}</version>
+</dependency>
+
+<!-- PipesParser and AsyncProcessor without PipesForkParser -->
+<dependency>
+ <groupId>org.apache.tika</groupId>
+ <artifactId>tika-pipes-core</artifactId>
+ <version>{tika-version}</version>
+</dependency>
+<dependency>
+ <groupId>org.apache.tika</groupId>
+ <artifactId>tika-parsers-standard-package</artifactId>
+ <version>{tika-version}</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.
+
+[source,java]
+----
+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 xref:using-tika/java-api/getting-started.adoc[Getting Started with the
Java API] and
+https://github.com/apache/tika/blob/main/tika-example/src/main/java/org/apache/tika/example/PipesForkParserExample.java[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 xref:pipes/configuration.adoc[Pipeline Configuration]
for the `pipes`
+block and xref:pipes/plugins/index.adoc[Plugins] for each fetcher and emitter.
+
+[source,json]
+----
+{
+ "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" }
+ }
+}
+----
+
+[source,java]
+----
+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:
+
+[source,java]
+----
+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:
+
+[source,json]
+----
+ "pipes-iterator": {
+ "file-system-pipes-iterator": {
+ "basePath": "/data/input",
+ "fetcherId": "fsf",
+ "emitterId": "fse"
+ }
+ },
+ "pipes-reporters": {
+ "file-system-reporter": { "statusFile": "/data/status.json" }
+ }
+----
+
+[source,java]
+----
+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
+xref:pipes/reporters.adoc[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:
+
+[cols="1,2,2"]
+|===
+|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`,
+`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 xref:pipes/cpu-sizing.adoc[Forked-JVM CPU and Heap
Sizing], and the
+timeout model in xref:pipes/timeouts.adoc[Timeouts].