Skip to content

Drain the stream queue across N consumer threads (#6069) - #6069

Open
FriedCosey wants to merge 3 commits into
pytorch:mainfrom
FriedCosey:export-D113599253
Open

Drain the stream queue across N consumer threads (#6069)#6069
FriedCosey wants to merge 3 commits into
pytorch:mainfrom
FriedCosey:export-D113599253

Conversation

@FriedCosey

@FriedCosey FriedCosey commented Jul 25, 2026

Copy link
Copy Markdown

Summary:

X-link: https://github.com/facebookresearch/FBGEMM/pull/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.

Reviewed By: chouxi

Differential Revision: D113599253

@meta-cla meta-cla Bot added the cla signed label Jul 25, 2026
@meta-codesync

meta-codesync Bot commented Jul 25, 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 D113599253.

@meta-codesync meta-codesync Bot changed the title Drain the stream queue across N consumer threads in raw_embedding_streamer Drain the stream queue across N consumer threads in raw_embedding_streamer (#6069) Jul 25, 2026
FriedCosey pushed a commit to FriedCosey/FBGEMM that referenced this pull request Jul 25, 2026
…eamer (pytorch#6069)

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
FriedCosey pushed a commit to FriedCosey/FBGEMM that referenced this pull request Jul 25, 2026
…eamer (pytorch#6069)

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
FriedCosey pushed a commit to FriedCosey/FBGEMM that referenced this pull request Jul 25, 2026
…eamer (pytorch#6069)

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
Joey Yang added 3 commits July 30, 2026 02:03
Summary:
The no-fbcode streamer test compiled the header with FBGEMM_FBCODE off but linked the always-flag-on :raw_embedding_streamer, so the flag-on constructor overran the smaller flag-off object layout → ASan heap-buffer-overflow (pre-existing latent bug).

Fix: link the test against a new -UFBGEMM_FBCODE build, :raw_embedding_streamer_no_fbcode, so its header view matches the library layout; switch the test include to the short form so it resolves via fbgemm_gpu's include/ dir instead of the manual pin to the flag-on lib.

Differential Revision: D114193888
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.

Reviewed By: chouxi

Differential Revision: D113599253
@meta-codesync meta-codesync Bot changed the title Drain the stream queue across N consumer threads in raw_embedding_streamer (#6069) Drain the stream queue across N consumer threads (#6069) Jul 30, 2026
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