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].

Reply via email to