Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions ENVIRONMENT.md
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@
| -------------------------------------- | -------- | -------------- | -------- | ---------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- |
| KAFKA_${prefix}BROKERS | string[] | localhost:9092 | | "prefix" is not required, can be used to have multiple kafka configurations: i.e: KAFKA_INSTANCE_A_BROKERS=kafka-a-01:9092,kafka-a-02:9092, KAFKA_INSTANCE_B_BROKERS=kafka-b-01:9092,kafka-b-02:9092 |
| KAFKA_${prefix}CLIENT_ID | string | _function_ | | default to `hostname()` |
| KAFKA_${prefix}LOG_LEVEL | string | NOTHING | | possible values: 'NOTHING', 'ERROR', 'WARN', 'INFO', 'DEBUG' |
| KAFKA_${prefix}SSL_CA_PATH | string | | | |
| KAFKA_${prefix}SSL_CERT_PEM_PATH | string | | | |
| KAFKA_${prefix}SSL_KEY_PEM_PATH | string | | | |
Expand Down
13 changes: 12 additions & 1 deletion src/interface.ts
Original file line number Diff line number Diff line change
@@ -1,14 +1,25 @@
import { EachMessageHandler, ProducerRecord } from 'kafkajs';
import {
ConsumerEvents,
EachMessageHandler,
InstrumentationEvent,
ProducerEvents,
ProducerRecord,
ValueOf
} from 'kafkajs';

export interface Consumer {
connect(): Promise<void>;
setCallback(func: EachMessageHandler): void;
disconnect(): Promise<boolean>;
subscribe(topics: string[], autoCommit: boolean, fromBeginning: boolean): Promise<void>;
on<T>(event: ValueOf<ConsumerEvents>, listener: (event: InstrumentationEvent<T>) => void): void;
off(event: ValueOf<ConsumerEvents>): void;
}

export interface Producer {
send(message: ProducerRecord): Promise<boolean>;
connect(): Promise<void>;
disconnect(): Promise<boolean>;
on<T>(event: ValueOf<ProducerEvents>, listener: (event: InstrumentationEvent<T>) => void): void;
off(event: ValueOf<ProducerEvents>): void;
}
28 changes: 27 additions & 1 deletion src/kafka/kafka-consumer.ts
Original file line number Diff line number Diff line change
@@ -1,4 +1,13 @@
import { Consumer, EachMessageHandler, EachMessagePayload, TopicPartitionOffsetAndMetadata } from 'kafkajs';
import {
Consumer,
ConsumerEvents,
EachMessageHandler,
EachMessagePayload,
InstrumentationEvent,
RemoveInstrumentationEventListener,
TopicPartitionOffsetAndMetadata,
ValueOf
} from 'kafkajs';
import { Consumer as GenericConsumer } from '../interface';
import { getLogger } from '@fluidware-it/saddlebag';

Expand All @@ -9,6 +18,8 @@ export class KafkaConsumer implements GenericConsumer {

private connected = false;

private eventsListener: Record<string, RemoveInstrumentationEventListener<ValueOf<ConsumerEvents>>[]> = {};

private callback: EachMessageHandler = () => {
throw new Error('EachMessagePayload must be override before connect() method call');
};
Expand Down Expand Up @@ -49,6 +60,21 @@ export class KafkaConsumer implements GenericConsumer {
this.logger.info('connected');
}

public on<T>(event: ValueOf<ConsumerEvents>, listener: (event: InstrumentationEvent<T>) => void): void {
const off = this.consumer.on(event, listener);
if (!this.eventsListener[event]) {
this.eventsListener[event] = [];
}
this.eventsListener[event].push(off);
}

public off(event: ValueOf<ConsumerEvents>) {
if (this.eventsListener[event]) {
this.eventsListener[event].forEach(off => off());
delete this.eventsListener[event];
}
}

public async disconnect(): Promise<boolean> {
if (!this.connected) return true;
try {
Expand Down
26 changes: 25 additions & 1 deletion src/kafka/kafka-producer.ts
Original file line number Diff line number Diff line change
@@ -1,4 +1,11 @@
import { Producer, ProducerRecord } from 'kafkajs';
import {
InstrumentationEvent,
Producer,
ProducerEvents,
ProducerRecord,
RemoveInstrumentationEventListener,
ValueOf
} from 'kafkajs';
import { Producer as GenericProducer } from '../interface';
import { getLogger } from '@fluidware-it/saddlebag';

Expand All @@ -9,6 +16,8 @@ export class KafkaProducer implements GenericProducer {

private connected = false;

private eventsListener: Record<string, RemoveInstrumentationEventListener<ValueOf<ProducerEvents>>[]> = {};

constructor(producer: Producer) {
this.producer = producer;
}
Expand All @@ -33,6 +42,21 @@ export class KafkaProducer implements GenericProducer {
}
}

public on<T>(event: ValueOf<ProducerEvents>, listener: (event: InstrumentationEvent<T>) => void): void {
const off = this.producer.on(event, listener);
if (!this.eventsListener[event]) {
this.eventsListener[event] = [];
}
this.eventsListener[event].push(off);
}

public off(event: ValueOf<ProducerEvents>): void {
if (this.eventsListener[event]) {
this.eventsListener[event].forEach(off => off());
delete this.eventsListener[event];
}
}

public async connect(): Promise<void> {
if (!this.connected) {
await this.producer.connect();
Expand Down