Skip to content

Kafka Crab JS API Reference

Inaiat Henrique edited this page Jan 9, 2025 · 1 revision

Kafka Crab JS API Reference

Core Types

KafkaClientConfig

interface KafkaClientConfig {
  brokers: string[];
  clientId?: string;
  logLevel?: string;
  brokerAddressFamily?: string;
  securityProtocol?: SecurityProtocol;
  configuration?: { [key: string]: string };
}

ConsumerConfiguration

interface ConsumerConfiguration {
  groupId: string;
  configuration?: { [key: string]: string };
  fecth_metadata_timeout?: number;
}

Consumer API

KafkaConsumer

Main consumer class for receiving messages from Kafka.

Constructor

constructor(config: KafkaClientConfig, consumerConfiguration: ConsumerConfiguration)

Methods

subscribe(topic: string | TopicPartitionConfig[]): Promise<void>
recv(): Promise<Message | null>
on_events(callback: (error: Error | undefined, event: KafkaEvent) => void): void
disconnect(): Promise<void>
pause(): Promise<void>
resume(): Promise<void>
seek(topic: string, partition: number, offset: OffsetModel, timeout?: number): Promise<void>
commit(topic: string, partition: number, offset: number, mode: CommitMode): Promise<void>
assignment(): Promise<TopicPartition[]>

Event Types

enum KafkaEventName {
  PreRebalance = "PreRebalance",
  PostRebalance = "PostRebalance",
  CommitCallback = "CommitCallback"
}

interface KafkaEventPayload {
  action?: string;
  tpl: TopicPartition[];
  error?: string;
}

interface KafkaEvent {
  name: KafkaEventName;
  payload: KafkaEventPayload;
}

Message Types

interface Message {
  topic: string;
  partition: number;
  offset: number;
  timestamp: number;
  headers?: { [key: string]: string };
  key?: Buffer;
  payload: Buffer;
}

interface TopicPartition {
  topic: string;
  partition: number;
  offset: number;
}

interface TopicPartitionConfig {
  topic: string;
  all_offsets?: OffsetModel;
  partition_offset?: PartitionOffset[];
}

interface OffsetModel {
  offset: number;
  position?: number;
}

type CommitMode = "Sync" | "Async";

Producer API

KafkaProducer

Main producer class for sending messages to Kafka.

Constructor

constructor(config: KafkaClientConfig)

Methods

send(options: ProducerSendOptions): Promise<ProduceResult[]>
flush(timeout?: number): Promise<void>

Producer Types

interface ProducerSendOptions {
  topic: string;
  messages: ProducerMessage[];
}

interface ProducerMessage {
  key?: string | Buffer;
  payload: string | Buffer;
  headers?: { [key: string]: string };
  partition?: number;
  timestamp?: number;
}

interface ProduceResult {
  topic: string;
  partition: number;
  offset?: number;
  error?: Error;
}

Security Types

enum SecurityProtocol {
  Plaintext = "plaintext",
  Ssl = "ssl",
  SaslPlaintext = "sasl_plaintext",
  SaslSsl = "sasl_ssl"
}

Configuration Reference

All configuration properties from librdkafka are supported through the configuration object in both KafkaClientConfig and ConsumerConfiguration.

See librdkafka Configuration Properties for the complete list of available options.

Common configurations include:

  • enable.auto.commit: Enable/disable automatic offset committing
  • auto.commit.interval.ms: Auto commit interval
  • auto.offset.reset: What to do when there is no initial offset
  • message.timeout.ms: Message timeout

Clone this wiki locally