///
import { Transform } from 'stream';
declare namespace CommitTransformStream {
interface CommitFunction {
(commits: {
topic: String;
partition: Number;
offset: Number;
}[]): Promise | unknown;
}
interface Option {
commitFunction: CommitFunction;
commitInterval: number;
}
}
/**
* @callback CommitFunction
* @param commits {Array.<{topic: String, partition: Number, offset: Number }>}
* @returns {Promise|*}
*/
/**
* @private
*/
declare class CommitTransformStream extends Transform {
private _bufferedOffset;
private _options;
private _forceCommitTimeout?;
private _currentCommitPromise;
private _isDestroyed;
/**
* @param options {Object}
* @param options.commitFunction {CommitFunction}
* @param options.commitInterval {Number} A positive integer that specifies a minimal duration (in milliseconds)
* between two offset commit request
*/
constructor(options: CommitTransformStream.Option);
_popBufferedOffset(): any[];
_performCommit(): Promise;
_setForceCommitTimeout(): void;
_transform(message: any, unused: any, callback: any): void;
_flush(callback: any): void;
_internalDestroy(e: any): void;
}
export default CommitTransformStream;