replace raw copy threads with a persistent copy_executor_ - #6079
Open
FriedCosey wants to merge 7 commits into
Open
replace raw copy threads with a persistent copy_executor_#6079FriedCosey wants to merge 7 commits into
FriedCosey wants to merge 7 commits into
Conversation
added 7 commits
July 27, 2026 12:44
Summary: X-link: facebookresearch/FBGEMM#2968 Speed up the RES C++ streamer's device->host copy by parallelizing it. Previously stream() copied all updated rows to CPU on a single thread and enqueued one item. This replaces that serial copy with a chunked copy across up to `kNumCopyThreads` threads: tensor_copy_chunk() copies a [start, end) row range, and the per-thread tiling is factored into a pure, unit-testable computeChunkRanges(). Differential Revision: D113594180
Summary: X-link: facebookresearch/FBGEMM#2971 Speed up the RES C++ streamer's ship stage by parallelizing it. Previously a single stream thread drained the queue, so the setEmbeddings RPCs to the PS ran one at a time. This diff replaces it with kNumConsumerThreads consumers on a folly::UMPMCQueue that ship concurrently. Differential Revision: D113599253
…#6075) Summary: X-link: meta-pytorch/torchrec#4464 X-link: facebookresearch/FBGEMM#2976 Make the RES streamer's chunk size and ship/copy thread counts tunable from config instead of compile-time constants. Previously kChunkSize / kNumConsumerThreads / kNumCopyThreads were constexpr in raw_embedding_streamer.cpp, so tuning them meant a base rebuild. This plumbs them as res_chunk_size / res_num_consumers / res_num_copy_threads from TableBatchedEmbeddingConfig through fused_params to the RawEmbeddingStreamer ctor, mirroring the existing res_store_shards. Defaults stay 500000/8/4, so behavior is unchanged until overridden Differential Revision: D113532922
Summary: X-link: facebookresearch/FBGEMM#2975 The four tensor copies in tensor_copy_chunk are independent, so dispatch each under its own sibling FBGEMM_DISPATCH_* instead of nesting weights -> indices -> identities -> runtime_meta. Two wins, no behavior change: 1. Readability: the flat structure drops the value_t / index_t / id_t / rm_t aliases that existed only to dodge scalar_t name-shadowing between the nested lambdas -- each copy now just uses scalar_t and reads top-to-bottom instead of 3 levels deep. 2. Fewer template instantiations: nesting stamps each inner copy once per outer type (multiplicative, ~4 x 2); siblings stamp each once per its own type set (additive, 4 + 2) -> smaller binary, faster compile. Differential Revision: D113533821
Summary: X-link: facebookresearch/FBGEMM#2978 Replace the per-iteration raw std::thread in the non-blocking stream() dispatch with a persistent named size-1 folly::CPUThreadPoolExecutor + SemiFuture. Kills per-iter thread create/join churn, names the thread in traces, and removes the raw-thread std::terminate risk (folly captures exceptions into the future) Differential Revision: D113669440
…orch#6076) Summary: X-link: facebookresearch/FBGEMM#2979 The consumer/ship path drained a hand-rolled folly::UMPMCQueue (weights_to_stream_queue_) with a pool of raw std::thread consumers that ran a 10ms-backoff polling loop try_dequeue(), sleep(10ms) on empty, retry and exited via a stop_ latch. This diff replaces both the queue and the raw thread pool with a persistent folly::CPUThreadPoolExecutor: each ship item is posted with add(...), and the dtor join()s the executor so all pending + in-flight ship tasks drain before teardown. Differential Revision: D113669447
Summary: X-link: facebookresearch/FBGEMM#2982 Swap the per-stream raw std::thread copy pool for a persistent folly::CPUThreadPoolExecutor; chunk copies fan out + join via collectAllRange; We still join so each copy finishes before the next iteration so this only removes per-stream thread churn. Differential Revision: D113752426
Contributor
|
@FriedCosey has exported this pull request. If you are a Meta employee, you can view the originating Diff in D113752426. |
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.
Summary:
X-link: https://github.com/facebookresearch/FBGEMM/pull/2982
Swap the per-stream raw std::thread copy pool for a persistent folly::CPUThreadPoolExecutor; chunk copies fan out + join via collectAllRange; We still join so each copy finishes before the next iteration so this only removes per-stream thread churn.
Differential Revision: D113752426