Conversation
The buffered pipelines keep one audio or feature bufferer per stream id and remove it only while processing a request with `is_last=True`. `delete_state()` removed only the stream's state, and `close_session()` and `open_session()` removed no bufferer, so a caller that ends its streams with `delete_state()`, as the in-tree simulstream adapter does, kept one bufferer per ended stream for the life of the pipeline. A new stream on the id of such a stream also continued the old stream's audio buffer. `BasePipeline.delete_state()` now removes the stream's bufferer, and `reset_session()` removes all of them, when the pipeline's bufferer is the buffered pipelines' `BatchedAudioBufferer` or `BatchedFeatureBufferer`, as a request with `is_last=True` does. `BufferedSALMPipeline`, which keeps its bufferer as `audio_bufferer`, does the same in its own `delete_state()` and `reset_session()`. Removing the bufferer of a stream that holds none, such as one whose `is_last` request already removed it or an unknown id, does nothing. The cache-aware pipelines' bufferer is left alone. Outputs of streams that end with `is_last` are unchanged: text, segments or words with their times and confidences, and finalization delays through `pipeline.run()` on 12 LibriSpeech test-other utterances, for the buffered RNN-T and CTC pipelines with frame and feature-buffer requests, on CPU. Part of NVIDIA-NeMo#16309 Signed-off-by: Zaheer Sheriff K <zaheersheriff.k@gmail.com>
4 of 7 tasks
This branch has not been deployed
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
What does this PR do ?
Makes
delete_state()in the buffered streaming pipelines (RNN-T, CTC and SALM) drop the stream's per-stream bufferer, andopen_session()/close_session()drop all of them. A stream that ends without anis_last=Truerequest no longer leaves its bufferer behind, and a later stream on the same id starts from an empty buffer.Collection: ASR
Part of #16309. #16321 frees the cache-aware pipelines' slots and leaves the buffered pipelines' per-stream bufferers out of scope; this PR covers those. It does not touch the cache-aware pipelines,
context_manager.pyorcache_feature_bufferer.py. #16321's description links the issue with a closing keyword, so the issue will be closed when #16321 merges; this PR is the buffered part.Changelog
BasePipeline.delete_state(): when the pipeline'sbuffereris aBatchedAudioBuffererorBatchedFeatureBufferer(whatinit_bufferer_for_buffered_streaming()builds forBufferedRNNTPipelineandBufferedCTCPipeline), also call its existingrm_bufferer(stream_id), the call theis_lastpath makes.BasePipeline.reset_session()(called byopen_session()andclose_session()): under the same condition, call the bufferer's existingreset(). Its docstring already reads "Reset the frame buffer and internal state pool".BufferedSALMPipeline: newdelete_state()andreset_session()overrides doing the same for itsaudio_bufferer(aBatchedIncrementalAudioBufferer), which it does not store asbufferer.tests/collections/asr/inference/test_buffered_bufferer_release.py: 5 tests, 23 cases.The problem
On
main(00278b0) each of the three buffered bufferers keeps a dict with one bufferer per stream id. It adds the entry on the stream's first request and removes it only while processing a request withis_last(audio_bufferer.pyL129-130,feature_bufferer.pyL158-159,incremental_audio_bufferer.pyL171-172). Onmain,BasePipeline.delete_state()only drops the state (base_pipeline.pyL141-144),reset_session()only clears the state pool (L158-160), and none of the buffered pipelines overrides either.The in-tree simulstream adapter ends every stream without
is_last:process_chunk()sendsis_last=False(simulstream_pipeline_adapter.pyL311) andend_of_stream()callsdelete_state()(L417). It forcesrequest_type='frame'(L137), so on these pipelines it usesBatchedAudioBufferer.This has two effects:
is_first, so a new stream on an id whose bufferer is still held continues that buffer: its first requests see the old stream's last audio as left context and a smaller left padding. A feature-buffer request carries its whole buffer, whichFeatureBufferer.update()copies in (feature_bufferer.pyL84), so there the leftover bufferer only costs memory.The adapter does not reuse ids within one instance (
clear()incrementsstream_id). But every instance starts atstream_id = 0(L96), and all instances share one class-level pipeline (load_model()returns early once it is built, L118-119). So in a process that creates a second adapter instance, its first stream starts on the id of the first instance's first stream. I do not know whether simulstream creates more than one instance per process; I measured what happens when it does.Measurements
All runs were on CPU in float32 with local checkpoints:
stt_en_conformer_transducer_smallinbuffered_rnnt.yamlandstt_en_conformer_ctc_smallinbuffered_ctc.yaml(both 4.8 s chunks, 8.0 s buffers). I overrode only model, device, dtype, ITN/NMT (off) and log level, plus the settings stated below.delete_state()and nois_lastrequest,mainheld 300 per-stream bufferers, and still 300 afterclose_session(). That held for RNN-T and CTC, through the adapter (frame requests) and throughtranscribe_step()directly (frame and feature-buffer requests). With this PR: 0.delete_state(). Onmainthe second stream's text changed in 3 of 12 pairs (RNN-T) and 4 of 12 (CTC), and its segments in 8 and 12. With this PR nothing changed. With feature-buffer requests nothing changed on either side.is_last.pipeline.run()outputs are identical betweenmainand this PR in all 8 configurations described below (16 runs).Bufferers held after 300 streams
One 4.8 s request of synthetic noise (
0.01 * N(0, 1)) per stream, 300 streams. The adapter runs use a stub for the two names the adapter imports from simulstream, which is not installed here; the adapter and pipeline code that ran is NeMo's own.main: after 300 streams / afterclose_session()process_chunk(),end_of_stream(),clear()), RNN-T and CTCtranscribe_step()thendelete_state(), RNN-T and CTCtranscribe_step()thendelete_state(), RNN-T and CTCtranscribe_step(), streams left open (RNN-T frame, CTC feature_buffer)is_last=True, RNN-T and CTCIn the
transcribe_step()runs onmain, the count also stayed at 300 after a followingopen_session(). The held buffer tensors alone were 512,000 bytes per stream with frame requests (153,600,000 for 300) and 256,320 bytes per stream with feature-buffer requests (76,896,000 for 300). I did not measure Python object overhead or process memory.I also reran the adapter harness from #16308 on buffered RNN-T: 300 streams, each with a biasing request. These are the counts after the 300th stream, and again after
close_session():mainStream-id reuse
I took 12 pairs of consecutive LibriSpeech test-other utterances (X, Y). For each pair I compared two runs of Y:
is_last=Falseand followed bydelete_state(7).The requests came from the pipeline's own request generator, as
pipeline.run()sends them. Y alone gave the same output twice in all 12 pairs.main: pairs where Y's text / segments (text, start, end, conf) differOn
main, start times changed in 8 RNN-T pairs, and confidences changed in all 12 CTC pairs. The text changed in 7 runs:An RNN-T example, where X (7902-96591-0008) ends with "... had a pretty good stock":
Through the adapter, I streamed the first 12 utterances in two ways: with one instance (stream ids not used before), and with a new instance per utterance (stream id 0 each time). On
mainthe transcripts of the two ways differed in 4 of 12 (RNN-T) and 5 of 12 (CTC). With this PR they differed in 0 of 12, and the new-instance transcripts equalmain's one-instance transcripts.Behaviour that should not change, and does not
I ran
pipeline.run()on the first 12 LibriSpeech test-other utterances withbatch_size=4, every stream ending with anis_lastrequest. The following are identical betweenmainand this PR for every stream:That covers RNN-T and CTC, frame and feature-buffer requests, and segment and word granularity: 8 configurations, 16 runs, 12 of 12 streams each. No bufferer was held after any of these runs on either side. On this path the new calls find nothing to remove: the
is_lastrequest has already removed the stream's bufferer whentranscribe_step()callsdelete_state()(base_pipeline.pyL304-305 in this PR, L299-300 onmain). On the merge of this PR with #16308 and #16321, the four segment-level configurations also matchmain.Why
BasePipeline, and composition with #16308, #16321 and #15827BasePipelinebuilds these two bufferers ininit_bufferer_for_buffered_streaming()(base_pipeline.pyL442 and L450 in this PR, L437 and L445 onmain), which theBufferedRNNTPipelineandBufferedCTCPipelineconstructors call. Itsreset_session()docstring already covers "the frame buffer". So I drop the bufferers inBasePipelineas well. Subclass overrides that end withsuper(), such as #16308'sBufferedRNNTPipeline.delete_state()and #15827'sBufferedRNNTPipeline.close_session(), reach the change without edits.The
isinstancecheck keepsBasePipelineoff every other bufferer, in particular the cache-awareBatchedCacheFeatureBufferer, which #16321 handles in its own overrides.test_base_pipeline_leaves_other_bufferers_alonepins that.BufferedSALMPipelinekeeps its bufferer asaudio_bufferer(buffered_salm_pipeline.pyL105), so it gets its own two small overrides. If you would rather have one pattern, for example overrides in each buffered pipeline or a small hook they implement, I am happy to change it.I merged this branch locally with each PR's current head (not pushed) and ran the PR test files and
tests/collections/asr/inference/:tests/collections/asr/inference/main#16308 and #16321 both add
CacheAwareRNNTPipeline.delete_state(), so they conflict with each other incache_aware_rnnt_pipeline.py, with or without this branch, which does not touch that file. Whichever of the two lands second needs a rebase. I resolved that conflict by keeping both bodies. The three-way merge then gives the same test results as #16308 + #16321 without this branch, plus this PR's 23 passing cases.The diffs of #16231, #16191 and #15940 (which change
base_pipeline.py) apply cleanly in either order with this change, as do #15186's hunks inbuffered_salm_pipeline.py. #15827, #16308 and #16321 do not touch the files changed here. #16158 targets another base branch. Itsbase_pipeline.pyhunks do not apply tomaineither, and they start at L211, away from the lines changed here.Tests
The tests are CPU-only and need no checkpoint. They build each buffered pipeline without a model. Its model step is replaced with the bufferer call that each real step begins with, recording the buffer and padding the model would get:
buffered_rnnt_pipeline.pyL416/L430,buffered_ctc_pipeline.pyL310/L324, andbuffered_salm_pipeline.pyL184 in this PR (L174 onmain).transcribe_step(),delete_state(), the sessions and the bufferers are the real ones. For RNN-T and CTC, the bufferers are the onesinit_bufferer_for_buffered_streaming()builds.test_delete_state_drops_bufferer_of_stream_ended_without_is_last(5 cases: RNN-T and CTC with frame and feature-buffer requests, and SALM): three streams dropped withdelete_state()leave no bufferer, and an unknown id is ignored.test_delete_state_drops_only_the_deleted_streams_bufferer(5 cases):delete_state()after anis_lastrequest does nothing. The stream still open keeps its bufferer and gets the same buffer on its next request as on a pipeline where nothing was deleted.test_close_and_open_session_drop_all_bufferers(5 cases).test_reused_stream_id_starts_from_an_empty_buffer(6 cases: the three frame configurations, with the stream ended bydelete_state()or byclose_session()+open_session()): a new stream on the old id sees the same buffer and padding as on a fresh pipeline. Feature-buffer requests are not included, since they carry their whole buffer.test_base_pipeline_leaves_other_bufferers_alone(2 cases, cache-aware RNN-T and CTC):BasePipeline.delete_state()andreset_session(), called directly, make no call on a bufferer of any other type. The test uses aMockthat records calls. It passes onmaintoo, sincemaincalls no bufferer there. Its job is to keep this change away from the cache-aware bufferer.The 7 skips need
--with_downloadsornemo_text_processing, and I did not run them.To check that the tests catch wrong fixes and not only a missing one, I swapped in eleven. Each failed some of the 23 cases. The two that loosen the
isinstancecheck are caught only bytest_base_pipeline_leaves_other_bufferers_alone.Wrong fixes and the cases they fail
delete_state()unchangeddelete_state()dropping every stream's buffererBatchedAudioBuffererself.buffererwithoutgetattrdelinstead ofrm_bufferer(), which raises for a stream that holds no buffererisinstancecheck replaced bybufferer is not Nonetest_base_pipeline_leaves_other_bufferers_alone)isinstancecheck replaced by ahasattrchecktest_base_pipeline_leaves_other_bufferers_alone)black==24.10.0andisort==5.13.2report no changes on the changed files.Limits
delete_state()is keyed by stream id only. A call that arrives after the id has been reused by a new stream deletes the new stream's state, as onmain, and now also its bufferer.pipeline.run()uses for feature-buffer requests has its ownBatchedAudioBufferer(multi_stream.pyL291), which also removes a stream only atis_last. It lives for onerun()call, anddelete_state()does not reach it. It is not changed here.Before your PR is "Ready for review"
Pre checks:
PR Type:
Additional Information
delete_state().cc @naymaraq @lilithgrigoryan @artbataev @nithinraok