Skip to content

basalt / queue-kafka/src / KafkaQueueDriver

Class: KafkaQueueDriver ​

Defined in: queue-kafka/src/index.ts:103

Kafka driver for @basaltkit/queue. Kafka is a log, not a task queue, so this driver is deliberately honest about what it can't do:

  • delayed and priority are NOT supported (Kafka has neither). With the queue's onUnsupported: 'throw' policy a delayed/priority dispatch fails loudly; with the default 'warn' it logs and proceeds without them.
  • retries are supported via a retry topic (<topic>.retry) the worker also consumes; exhausted jobs go to <topic>.dead. There is no backoff delay (Kafka can't defer a message), so backoff is not supported either.

Worker concurrency is bounded by the topic's partition count, not the concurrency number (passed through as partitionsConsumedConcurrently).

Implements ​

Constructors ​

Constructor ​

> new KafkaQueueDriver(options): KafkaQueueDriver

Defined in: queue-kafka/src/index.ts:122

Parameters ​

options ​

KafkaDriverOptions

Returns ​

KafkaQueueDriver

Properties ​

capabilities ​

> readonly capabilities: DriverCapabilities

Defined in: queue-kafka/src/index.ts:105

What this backend honors — see DriverCapabilities.

Implementation of ​

QueueDriver.capabilities


name ​

> readonly name: "kafka" = 'kafka'

Defined in: queue-kafka/src/index.ts:104

Short identifier used in diagnostics (e.g. 'bullmq', 'sync').

Implementation of ​

QueueDriver.name

Methods ​

add() ​

> add(queue, jobName, data, options): Promise&lt;void&gt;

Defined in: queue-kafka/src/index.ts:139

Parameters ​

queue ​

string

jobName ​

string

data ​

unknown

options ​

AddJobOptions

Returns ​

Promise&lt;void&gt;

Implementation of ​

QueueDriver.add


close() ​

> close(): Promise&lt;void&gt;

Defined in: queue-kafka/src/index.ts:178

Returns ​

Promise&lt;void&gt;

Implementation of ​

QueueDriver.close


setExecutor() ​

> setExecutor(executor): void

Defined in: queue-kafka/src/index.ts:135

Called once by the QueueManager — how to execute a received job.

Parameters ​

executor ​

JobExecutor

Returns ​

void

Implementation of ​

QueueDriver.setExecutor


startWorker() ​

> startWorker(queue, options?): void

Defined in: queue-kafka/src/index.ts:159

Starts a worker for the queue (no-op in the sync driver: add executes inline).

Parameters ​

queue ​

string

options? ​
concurrency? ​

number

Returns ​

void

Implementation of ​

QueueDriver.startWorker

Released under the MIT License.