/// import { Transform } from 'stream'; /** * @callback MessageConsumerCallback * @param message {Object} * @param message.topic {String} * @param message.offset {Number} * @param message.values {String|Buffer} * @param message.partition {Number} * @returns {Promise|*} 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. */ /** * @callback FailedMessageConsumerCallback * @param error * @param message {Object} * @param message.topic {String} * @param message.offset {Number} * @param message.values {String|Buffer} * @param message.partition {Number} * @returns {Promise|*} 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. * */ declare namespace ConsumeTransformStream { interface Option { consumeConcurrency: number; consumeTimeout: number; groupId: string; messageConsumer(message: any): Promise | unknown; failedMessageConsumer?(error: any, message: any): Promise | unknown; } } /** * @private */ declare class ConsumeTransformStream extends Transform { private _options; private _currentConsumeConcurrency; private _concurrentPromise; private _waitingQueue; private _lastMessageQueuedPromise; private _isDestroyed; private _unhandledException?; /** * * @param options.messageConsumer {MessageConsumerCallback} * @param options {Object} * @param options.consumeConcurrency {Number} * @param options.consumeTimeout {Number} * @param options.groupId {String} * @param [options.failedMessageConsumer] {FailedMessageConsumerCallback} */ constructor(options: ConsumeTransformStream.Option); _transform(message: any, encoding: any, callback: any): void; _flush(callback: any): void; _consumeMessage(message: any): any; _concurrentConsumeMessage(message: any): void; _enqueue(message: any): any; _dequeue(): void; _internalDestroy(e: Error): void; } export default ConsumeTransformStream;