///
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;