Skip to content

replace raw copy threads with a persistent copy_executor_ - #6079

Open
FriedCosey wants to merge 7 commits into
pytorch:mainfrom
FriedCosey:export-D113752426
Open

replace raw copy threads with a persistent copy_executor_#6079
FriedCosey wants to merge 7 commits into
pytorch:mainfrom
FriedCosey:export-D113752426

Conversation

@FriedCosey

Copy link
Copy Markdown

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

Joey Yang 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
@meta-cla meta-cla Bot added the cla signed label Jul 27, 2026
@meta-codesync

meta-codesync Bot commented Jul 27, 2026

Copy link
Copy Markdown
Contributor

@FriedCosey has exported this pull request. If you are a Meta employee, you can view the originating Diff in D113752426.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant