Skip to main content

API Reference

The exported surface of @nest-native/kafka. Everything below is imported from the package root.

Module​

ExportKindNotes
KafkaModuleclassforRoot(options), forRootAsync(options), forFeature([HandlerClass]) — see Module.
KafkaModuleOptionsinterfaceOptions for forRoot.
KafkaModuleAsyncOptionsinterfaceOptions for forRootAsync.
KafkaConcurrencyOptionsinterfaceconcurrency and maxInFlight, shared by module / consumer / handler.

Consumer Decorators​

ExportKindNotes
KafkaConsumerdecoratorClass-level: @KafkaConsumer(topic?, options?).
KafkaHandlerdecoratorMethod-level: @KafkaHandler(topic?, options?).
KafkaTopicPatterntypestring | RegExp — the topic both decorators take; a pattern is anchored, flag-free, POSIX syntax. See Topic Patterns.
KafkaConsumerOptionsinterfacegroupId plus concurrency options.
KafkaHandlerOptionsinterfacebatch, reply, plus concurrency options.

See Consumers.

Parameter Decorators​

ExportKindNotes
KafkaMessagedecoratorWhole payload, or one property with @KafkaMessage('prop').
KafkaHeadersdecoratorAll headers, or one with @KafkaHeaders('key').
KafkaCtxdecoratorThe raw KafkaContext.
KafkaBatchdecoratorThe raw KafkaConsumerBatch (batch handlers).
KafkaContextclassTransport context: getTopic(), getPartition(), getMessage(), getHeaders().
KafkaBatchContextclassBatch transport context.
KafkaMessageHeadersinterfaceThe header map shape.

See Parameter Decorators.

Producer​

ExportKindNotes
KafkaProducerServiceclasssend, sendBatch, transactional.
InjectKafkaProducerdecoratorInject the raw KafkaDriverProducer.
KafkaTransactioninterfaceThe transaction handle passed to transactional.
KafkaSendRecordinterfaceA single send payload.
KafkaSendBatchinterfaceA sendBatch payload.
KafkaTransactionOffsetsinterfaceThe sendOffsets argument (Confluent shape).

See Producer and Transactions.

Error Mapping​

ExportKindNotes
KafkaErrorMappertype(error, context) => KafkaErrorBehavior | Promise<KafkaErrorBehavior> — awaited; a rejection retries.
KafkaErrorBehaviortype'commit' | 'retry'.
defaultKafkaErrorMapperconstCommits 4xx client errors, retries the rest.
KafkaErrorContexttypeKafkaContext | KafkaBatchContext.
KafkaRetryBackoffOptionsinterfaceinitialDelayMs, maxDelayMs, multiplier for the retryBackoff module option.
DEFAULT_KAFKA_RETRY_BACKOFFconst{ initialDelayMs: 1000, maxDelayMs: 30000, multiplier: 2 }.
toDeadLetterMessagefunction(context, error, options?) → the dead-letter record: original key, value, headers, plus Spring's kafka_dlt-* headers.
toDeadLetterMessagesfunctionThe same for every message of a failed batch (KafkaBatchContext).
readDeadLetterHeadersfunctionDecodes the kafka_dlt-* headers of a consumed record; undefined when it is not a dead letter.
KAFKA_DEAD_LETTER_HEADERSconstThe header names, Spring Kafka's.
KafkaDeadLetterOptionsinterfaceconsumerGroup, includeStackTrace (default true).
KafkaDeadLetterInfointerfaceWhat readDeadLetterHeaders returns.

See Error Mapping.

Request-Reply​

Opt-in, and inert until requestReply is configured. See Request-Reply.

ExportKindNotes
KafkaRequestReplyServiceclassrequest(record, options?) — produce a request and await the correlated reply.
KafkaRequestReplyOptionsinterfaceThe requestReply module block: replyTopic (required), timeoutMs, readinessTimeoutMs, groupIdPrefix, headers, consumer.
KafkaRequestRecordinterface{ topic, message } — the request to produce.
KafkaRequestOptionsinterfacePer-call timeoutMs and signal.
KafkaReplyinterface{ value, headers, correlationId, topic, partition, offset? }.
KafkaRequestReplyHeaderKeysinterfaceThe five header names; defaults interoperate with @nestjs/microservices.
KafkaReplyTimeoutErrorclassNo reply in time — the outcome is unknown.
KafkaReplyRemoteErrorclassThe remote handler failed and its error mapped to 'commit'.
KafkaReplyAbortedErrorclassThe wait was cancelled by a signal or by shutdown.
KafkaReplyDeliveryErrorclassRaised on the replying side when the reply could not be produced.
DEFAULT_KAFKA_REQUEST_REPLY_HEADERSconstThe default header key map.
DEFAULT_REQUEST_TIMEOUT_MSconst30000.
DEFAULT_READINESS_TIMEOUT_MSconst10000.

Driver​

ExportKindNotes
createConfluentDriverconstThe default driver factory; lazily loads the Confluent client.
KafkaClientDriverinterfaceThe driver contract.
KafkaDriverProducerinterfaceThe producer the driver exposes.
KafkaDriverConsumerinterfaceThe 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.
KafkaDriverAdmininterfaceThe admin client the driver may expose (createAdmin) for the health indicator: connect, disconnect, listTopics.
KafkaTopicPartitionsinterfaceA topic and, optionally, the partitions to pause.
KafkaDriverFactorytypedriverFactory option shape.

The driver is an advanced seam. Most applications never touch it directly.

Health​

ExportKindNotes
KafkaHealthIndicatorinjectableisHealthy(key?, options?) — a metadata round trip; resolves {[key]: {status: 'up' | 'down', …}}, terminus's result shape.
KafkaHealthCheckOptionsinterfacetimeoutMs (default 5000).
KafkaHealthIndicatorResulttypeRecord<string, KafkaHealthStatus>.
KafkaHealthStatusinterfacestatus plus details (latencyMs, topics, or message).
DEFAULT_KAFKA_HEALTH_TIMEOUT_MSconst5000.

See Resilience.

Testing​

ExportKindNotes
KafkaTestModuleclassIn-memory transport: forRoot, forRootAsync.
InMemoryKafkaBrokerclassThe loopback broker: emit, idle, getSent, getSentTo.
KAFKA_TEST_BROKERsymbolInjection token for the broker.
InjectKafkaTestBrokerdecoratorInject the broker.
createMockKafkaProducerfunctionRecording producer mock.
createMockTransactionfunctionRecording 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.