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

kafka-emitter

KafkaEmitter

Iterator

kafka-pipes-iterator

KafkaPipesIterator

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

topic

required

Kafka topic to publish to (validated non-blank).

bootstrapServers

required

Comma-separated host:port list of Kafka brokers (validated non-blank).

acks

all

Producer acks: 0, 1, or all.

lingerMs

5

Producer linger, ms.

batchSize

16384

Producer batch size, bytes.

bufferMemory

33554432

Producer buffer memory, bytes (32 MiB).

compressionType

none

One of none, gzip, snappy, lz4, zstd.

connectionsMaxIdleMs

540000

Passed through to the producer as connections.max.idle.ms.

deliveryTimeoutMs

120000

End-to-end delivery timeout. Must be at least lingerMs + requestTimeoutMs.

enableIdempotence

true

Idempotent producer. Requires acks=all and maxInFlightRequestsPerConnection of at most 5.

interceptorClasses

none

Comma-separated list of producer interceptor class names.

maxBlockMs

60000

How long send() blocks when the buffer is full.

maxInFlightRequestsPerConnection

5

In-flight requests per connection. Must be at least 1; 0 fails producer construction.

maxRequestSize

1048576

Maximum request size, bytes (1 MiB).

metadataMaxAgeMs

300000

Metadata refresh interval.

requestTimeoutMs

30000

Request timeout.

retries

2147483647

Producer retries, capped in practice by deliveryTimeoutMs.

retryBackoffMs

100

Backoff between retries.

transactionTimeoutMs

60000

Transaction timeout; only meaningful with transactionalId.

transactionalId

none

Set to enable the transactional producer.

clientId

none

client.id sent with each request.

keySerializer / valueSerializer

StringSerializer

Fully-qualified serializer class names. Unset (or unloadable) falls back to StringSerializer for both.

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

topic

required

Kafka topic to consume from.

bootstrapServers

required

Broker list.

groupId

none

Kafka consumer group ID. Strongly recommended in production for failover and partition reassignment.

keySerializer / valueSerializer

StringDeserializer

Consumer deserializer class names. Unset (or unloadable) falls back to StringDeserializer for both.

autoOffsetReset

earliest

What to do on first connect: earliest or latest.

pollDelayMs

100

Timeout passed to each poll() call.

emitMax

-1

Maximum tuples to emit. -1 means unbounded.

assignmentTimeoutMs

30000

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.

drainIdleMs

1000

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 pollDelayMs is.

groupInitialRebalanceDelayMs

3000

Deprecated and ignored. This is a broker setting, not a consumer one, so Kafka never applied it. Use assignmentTimeoutMs instead. Still accepted so existing configs start; scheduled for removal.

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-clients SDK.

  • The emitter is fire-and-forget at the Tika level; durability is determined by Kafka’s acks and broker replication factor, not by Tika.

  • For exactly-once semantics, set enableIdempotence: true (and ensure acks: all); for transactional semantics, also set transactionalId.

  • The iterator’s groupId controls partition assignment. Set it explicitly — without one, the consumer receives a transient assignment that resets on restart.