-
Notifications
You must be signed in to change notification settings - Fork 0
Kafka Crab JS API Reference
Inaiat Henrique edited this page Jan 9, 2025
·
1 revision
interface KafkaClientConfig {
brokers: string[];
clientId?: string;
logLevel?: string;
brokerAddressFamily?: string;
securityProtocol?: SecurityProtocol;
configuration?: { [key: string]: string };
}interface ConsumerConfiguration {
groupId: string;
configuration?: { [key: string]: string };
fecth_metadata_timeout?: number;
}Main consumer class for receiving messages from Kafka.
constructor(config: KafkaClientConfig, consumerConfiguration: ConsumerConfiguration)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[]>enum KafkaEventName {
PreRebalance = "PreRebalance",
PostRebalance = "PostRebalance",
CommitCallback = "CommitCallback"
}
interface KafkaEventPayload {
action?: string;
tpl: TopicPartition[];
error?: string;
}
interface KafkaEvent {
name: KafkaEventName;
payload: KafkaEventPayload;
}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";Main producer class for sending messages to Kafka.
constructor(config: KafkaClientConfig)send(options: ProducerSendOptions): Promise<ProduceResult[]>
flush(timeout?: number): Promise<void>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;
}enum SecurityProtocol {
Plaintext = "plaintext",
Ssl = "ssl",
SaslPlaintext = "sasl_plaintext",
SaslSsl = "sasl_ssl"
}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