Apache Kafka Plugin
The Apache Kafka plugin (tika-pipes-kafka) provides an emitter (publishes parsed documents to a Kafka topic) and an iterator (consumes fetch requests from a Kafka topic).
The two halves have different standing. The emitter is a plain producer and a good fit: Tika parses, results stream to a topic for downstream indexing. The iterator is a worked example whose suitability depends on your parse latency — see Things to know before building on it before building on it.
| Interface | Component name | Class |
|---|---|---|
Emitter |
|
|
Iterator |
|
|
Kafka Emitter (kafka-emitter)
Publishes each parsed document as a record to a Kafka topic.
{
"emitters": {
"kafe": {
"kafka-emitter": {
"topic": "tika-parsed-docs",
"bootstrapServers": "kafka1.example.com:9092,kafka2.example.com:9092",
"acks": "all",
"lingerMs": 5000,
"batchSize": 16384,
"compressionType": "lz4",
"enableIdempotence": true,
"maxRequestSize": 1048576,
"requestTimeoutMs": 30000,
"deliveryTimeoutMs": 120000,
"clientId": "tika-pipes-emitter"
}
}
}
}
Configuration
Fields map onto standard Kafka producer settings. Tika supplies no values of its own: an absent field is omitted from the producer properties, so the Kafka client’s own default applies. The Default column below is therefore the Kafka client’s default.
| Field | Default | Description |
|---|---|---|
|
required |
Kafka topic to publish to (validated non-blank). |
|
required |
Comma-separated |
|
|
Producer acks: |
|
|
Producer linger, ms. |
|
|
Producer batch size, bytes. |
|
|
Producer buffer memory, bytes (32 MiB). |
|
|
One of |
|
|
Passed through to the producer as |
|
|
End-to-end delivery timeout. Must be at least |
|
|
Idempotent producer. Requires |
|
none |
Comma-separated list of producer interceptor class names. |
|
|
How long |
|
|
In-flight requests per connection. Must be at least |
|
|
Maximum request size, bytes (1 MiB). |
|
|
Metadata refresh interval. |
|
|
Request timeout. |
|
|
Producer retries, capped in practice by |
|
|
Backoff between retries. |
|
|
Transaction timeout; only meaningful with |
|
none |
Set to enable the transactional producer. |
|
none |
|
|
|
Fully-qualified serializer class names. Unset (or unloadable) falls back to |
Kafka Iterator (kafka-pipes-iterator)
|
Reference example — check that your parse latency suits it
Kafka’s consumer model assumes bounded, roughly uniform per-message processing time, so how well this component works depends on how long your documents take to parse. Short, predictable parses fit that assumption. Long-running parses do not: pairing Kafka with documents that take minutes to an hour — OCR’d PDFs, large archives, anything near the one-hour default total-task timeout — works against the grain of the offset model, and the caveats below stop being theoretical. It is kept as a worked example rather than a supported production integration. If your workload has a long latency tail, consider driving tika-server or tika-grpc from your own consumer, where you control acknowledgement and retry, and use the Kafka emitter to publish results. |
Things to know before building on it
Head-of-line blocking, in proportion to your latency spread. A partition is consumed in order
by one member of the group. Even with numClients workers parsing in parallel, correct offset
handling can only advance the commit watermark to the lowest un-acknowledged offset, so one slow
document holds up its partition’s progress. With short, uniform parses this is barely
noticeable; with a long tail, one document can stall a partition for as long as it parses.
This follows from Kafka’s offset model rather than from a setting.
At-most-once delivery. enable.auto.commit is left at Kafka’s default of true, so offsets
are committed on a timer as soon as records are polled — before Tika has parsed or emitted
them. A crash, OOM or failed emit in between loses those documents silently. Committing after a
successful emit is not currently possible: the iterator pushes tuples onto a queue and receives
no completion signal back.
It drains and exits; it does not stream. The iterator enqueues what is on the topic and then
finishes, because tika-pipes' iterator contract is finite. A quiet period ends the run (see
drainIdleMs). It suits a periodic batch drain, not a long-running consumer.
Give the consumer room to join. A stock Kafka broker applies
group.initial.rebalance.delay.ms (default 3000) to the first member joining an empty group.
A deployment that runs back-to-back, or keeps another member in the group, never pays this; one
that starts cold pays it every run. The iterator waits for a partition assignment up to
assignmentTimeoutMs rather than mistaking a not-yet-assigned consumer for an empty topic.
Configuration
In addition to the required fetcherId / emitterId (see Wiring an Iterator):
| Field | Default | Description |
|---|---|---|
|
required |
Kafka topic to consume from. |
|
required |
Broker list. |
|
none |
Kafka consumer group ID. Strongly recommended in production for failover and partition reassignment. |
|
|
Consumer deserializer class names. Unset (or unloadable) falls back to |
|
|
What to do on first connect: |
|
|
Timeout passed to each |
|
|
Maximum tuples to emit. |
|
|
How long to wait for the consumer to be assigned a partition before failing. A newly subscribed consumer returns empty polls while it joins the group; the iterator waits for an assignment so it cannot mistake that for an empty topic. |
|
|
How long the topic must stay quiet (no records, after assignment) before it is treated as
drained and the iterator finishes. A duration rather than a poll count, so it holds however
short |
|
|
Deprecated and ignored. This is a broker setting, not a consumer one, so Kafka never
applied it. Use |
Complete Pipeline Example
A Kafka iterator (consuming fetch requests), a filesystem fetcher, and a Kafka emitter (publishing parsed results). This end-to-end shape is the reference example; see Things to know before building on it before relying on the iterator half.
{
"content-handler-factory": {
"basic-content-handler-factory": {
"type": "TEXT",
"writeLimit": -1,
"throwOnWriteLimitReached": true
}
},
"fetchers": {
"fsf": {
"file-system-fetcher": {
"basePath": "/data/input",
"extractFileSystemMetadata": false
}
}
},
"emitters": {
"kafe": {
"kafka-emitter": {
"topic": "tika-parsed-docs",
"bootstrapServers": "kafka1.example.com:9092",
"acks": "all",
"compressionType": "lz4",
"enableIdempotence": true
}
}
},
"pipes-iterator": {
"kafka-pipes-iterator": {
"topic": "tika-fetch-requests",
"bootstrapServers": "kafka1.example.com:9092",
"groupId": "tika-pipes-iterator",
"autoOffsetReset": "earliest",
"fetcherId": "fsf",
"emitterId": "kafe"
}
},
"pipes": {
"parseMode": "RMETA",
"onParseException": "EMIT",
"numClients": 4
}
}
Notes
-
The Kafka plugin uses the official
kafka-clientsSDK. -
The emitter is fire-and-forget at the Tika level; durability is determined by Kafka’s
acksand broker replication factor, not by Tika. -
For exactly-once semantics, set
enableIdempotence: true(and ensureacks: all); for transactional semantics, also settransactionalId. -
The iterator’s
groupIdcontrols partition assignment. Set it explicitly — without one, the consumer receives a transient assignment that resets on restart.