/// import { Transform } from 'stream'; import { Message } from 'kafka-node'; export interface MessageConsumer { /** * @param message - Message to be consumed * @returns This function need to return a promise if the message is consumed asynchronously, * value of any other types indicates that the message have already been consumed. */ (message: Message): Promise | unknown; } export interface FailedMessageConsumer { /** * @param error - The error raised while consuming the message * @param message - The message failed to be consumed * * This function need to return a promise if the message is consumed asynchronously, * value of any other types indicates that the message have already been consumed. */ (error: Error, message: Message): Promise | unknown; } /** * ConsumeOption */ export interface ConsumeOption { /** * How many message could be consumed concurrently per partition */ consumeConcurrency: number; /** * Timeout of consuming procedure for a single message */ consumeTimeout: number; /** * The group that the consumer is belonging to */ groupId: string; /** * The consuming procedure to be invoked for each message */ messageConsumer: MessageConsumer; /** * Optional failed message handler */ failedMessageConsumer?: FailedMessageConsumer; } /** * The consuming part of the pipeline * * @private */ declare class ConsumeStream extends Transform { private _options; private _currentConsumeConcurrency; private _concurrentPromise; private _waitingQueue; private _lastMessageQueuedPromise; private _isDestroyed; private _unhandledException?; /** * @param options - Option controls the behaviors that how the message is consumed */ constructor(options: ConsumeOption); /** * * @param message * @param encoding * @param callback * @private */ _transform(message: Message, encoding: any, callback: any): void; _flush(callback: any): void; /** * * @param message * @private */ private _consumeMessage; private _concurrentConsumeMessage; private _enqueue; private _dequeue; private _internalDestroy; } export default ConsumeStream;