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