API Reference
The exported surface of @nest-native/kafka. Everything below is imported from
the package root.
Module
| Export | Kind | Notes |
|---|---|---|
KafkaModule | class | forRoot(options), forRootAsync(options), forFeature([HandlerClass]) — see Module. |
KafkaModuleOptions | interface | Options for forRoot. |
KafkaModuleAsyncOptions | interface | Options for forRootAsync. |
KafkaConcurrencyOptions | interface | concurrency and maxInFlight, shared by module / consumer / handler. |
Consumer Decorators
| Export | Kind | Notes |
|---|---|---|
KafkaConsumer | decorator | Class-level: @KafkaConsumer(topic?, options?). |
KafkaHandler | decorator | Method-level: @KafkaHandler(topic?, options?). |
KafkaTopicPattern | type | string | RegExp — the topic both decorators take; a pattern is anchored, flag-free, POSIX syntax. See Topic Patterns. |
KafkaConsumerOptions | interface | groupId plus concurrency options. |
KafkaHandlerOptions | interface | batch, reply, plus concurrency options. |
See Consumers.
Parameter Decorators
| Export | Kind | Notes |
|---|---|---|
KafkaMessage | decorator | Whole payload, or one property with @KafkaMessage('prop'). |
KafkaHeaders | decorator | All headers, or one with @KafkaHeaders('key'). |
KafkaCtx | decorator | The raw KafkaContext. |
KafkaBatch | decorator | The raw KafkaConsumerBatch (batch handlers). |
KafkaContext | class | Transport context: getTopic(), getPartition(), getMessage(), getHeaders(). |
KafkaBatchContext | class | Batch transport context. |
KafkaMessageHeaders | interface | The header map shape. |
See Parameter Decorators.
Producer
| Export | Kind | Notes |
|---|---|---|
KafkaProducerService | class | send, sendBatch, transactional. |
InjectKafkaProducer | decorator | Inject the raw KafkaDriverProducer. |
KafkaTransaction | interface | The transaction handle passed to transactional. |
KafkaSendRecord | interface | A single send payload. |
KafkaSendBatch | interface | A sendBatch payload. |
KafkaTransactionOffsets | interface | The sendOffsets argument (Confluent shape). |
See Producer and Transactions.
Error Mapping
| Export | Kind | Notes |
|---|---|---|
KafkaErrorMapper | type | (error, context) => KafkaErrorBehavior | Promise<KafkaErrorBehavior> — awaited; a rejection retries. |
KafkaErrorBehavior | type | 'commit' | 'retry'. |
defaultKafkaErrorMapper | const | Commits 4xx client errors, retries the rest. |
KafkaErrorContext | type | KafkaContext | KafkaBatchContext. |
KafkaRetryBackoffOptions | interface | initialDelayMs, maxDelayMs, multiplier for the retryBackoff module option. |
DEFAULT_KAFKA_RETRY_BACKOFF | const | { initialDelayMs: 1000, maxDelayMs: 30000, multiplier: 2 }. |
toDeadLetterMessage | function | (context, error, options?) → the dead-letter record: original key, value, headers, plus Spring's kafka_dlt-* headers. |
toDeadLetterMessages | function | The same for every message of a failed batch (KafkaBatchContext). |
readDeadLetterHeaders | function | Decodes the kafka_dlt-* headers of a consumed record; undefined when it is not a dead letter. |
KAFKA_DEAD_LETTER_HEADERS | const | The header names, Spring Kafka's. |
KafkaDeadLetterOptions | interface | consumerGroup, includeStackTrace (default true). |
KafkaDeadLetterInfo | interface | What readDeadLetterHeaders returns. |
See Error Mapping.
Request-Reply
Opt-in, and inert until requestReply is configured. See
Request-Reply.
| Export | Kind | Notes |
|---|---|---|
KafkaRequestReplyService | class | request(record, options?) — produce a request and await the correlated reply. |
KafkaRequestReplyOptions | interface | The requestReply module block: replyTopic (required), timeoutMs, readinessTimeoutMs, groupIdPrefix, headers, consumer. |
KafkaRequestRecord | interface | { topic, message } — the request to produce. |
KafkaRequestOptions | interface | Per-call timeoutMs and signal. |
KafkaReply | interface | { value, headers, correlationId, topic, partition, offset? }. |
KafkaRequestReplyHeaderKeys | interface | The five header names; defaults interoperate with @nestjs/microservices. |
KafkaReplyTimeoutError | class | No reply in time — the outcome is unknown. |
KafkaReplyRemoteError | class | The remote handler failed and its error mapped to 'commit'. |
KafkaReplyAbortedError | class | The wait was cancelled by a signal or by shutdown. |
KafkaReplyDeliveryError | class | Raised on the replying side when the reply could not be produced. |
DEFAULT_KAFKA_REQUEST_REPLY_HEADERS | const | The default header key map. |
DEFAULT_REQUEST_TIMEOUT_MS | const | 30000. |
DEFAULT_READINESS_TIMEOUT_MS | const | 10000. |
Driver
| Export | Kind | Notes |
|---|---|---|
createConfluentDriver | const | The default driver factory; lazily loads the Confluent client. |
KafkaClientDriver | interface | The driver contract. |
KafkaDriverProducer | interface | The producer the driver exposes. |
KafkaDriverConsumer | interface | The consumer the driver exposes. An optional pause lets graceful shutdown stop deliveries before draining; a custom driver without it loses nothing, since late records are handed back. |
KafkaDriverAdmin | interface | The admin client the driver may expose (createAdmin) for the health indicator: connect, disconnect, listTopics. |
KafkaTopicPartitions | interface | A topic and, optionally, the partitions to pause. |
KafkaDriverFactory | type | driverFactory option shape. |
The driver is an advanced seam. Most applications never touch it directly.
Health
| Export | Kind | Notes |
|---|---|---|
KafkaHealthIndicator | injectable | isHealthy(key?, options?) — a metadata round trip; resolves {[key]: {status: 'up' | 'down', …}}, terminus's result shape. |
KafkaHealthCheckOptions | interface | timeoutMs (default 5000). |
KafkaHealthIndicatorResult | type | Record<string, KafkaHealthStatus>. |
KafkaHealthStatus | interface | status plus details (latencyMs, topics, or message). |
DEFAULT_KAFKA_HEALTH_TIMEOUT_MS | const | 5000. |
See Resilience.
Testing
| Export | Kind | Notes |
|---|---|---|
KafkaTestModule | class | In-memory transport: forRoot, forRootAsync. |
InMemoryKafkaBroker | class | The loopback broker: emit, idle, getSent, getSentTo. |
KAFKA_TEST_BROKER | symbol | Injection token for the broker. |
InjectKafkaTestBroker | decorator | Inject the broker. |
createMockKafkaProducer | function | Recording producer mock. |
createMockTransaction | function | Recording transaction mock. |
See Testing. Import the testing utilities from the
@nest-native/kafka/testing entrypoint — they are kept out of the package root
so test scaffolding never enters a consumer's production import surface.