Repository navigation
[Core][Frontend] Support request-level batching for diffusion pipelines - #4079
Conversation
980da0e to
4f3e5ca
Compare
|
Codex usage limits have been reached for code reviews. Please check with the admins of this repo to increase the limits by adding credits. |
7771f3c to
c1859ab
Compare
| allow for future improvements. | ||
| Diffusion list-prompt input is not represented as one packed | ||
| `OmniDiffusionRequest`. Submit multiple prompts as independent requests to | ||
| use automatic request-level batching. |
There was a problem hiding this comment.
how to control the batch size?
There was a problem hiding this comment.
Use max_num_seqs, consistent with the step-wise batching path
|
I think it's better to compare step-level batching and request-level batching on time performance |
There was a problem hiding this comment.
This refactor introduces request-level batching but is missing some critical validation and architectural cleanup.
vllm_omni/diffusion/request.py:35: If diffusion models later support multi-part single prompts (e.g., interleaved text and image dicts in a list), this check (isinstance(self.prompt, list)) will incorrectly reject them. Consider checking if the elements are valid single-prompt types rather than rejecting all lists.vllm_omni/diffusion/worker/diffusion_model_runner.py:310: Hardcodingtea_cachespecific logic in the model runner breaks encapsulation. The cache backend interface should define whether it requiresnum_inference_steps, or gracefully handleNonewithout runner-side workarounds.vllm_omni/entrypoints/async_omni.py:313: Since this is a user-facing API breaking change, ensure that this error message also points users to the updated documentation or an example script for the new independent request submission pattern.
|
seems it does not support batching variable-length prompts in my test |
|
@knlnguyen1802 PTAL |
any updates? |
… structure across various models Signed-off-by: jader <yjader@foxmail.com>
Currently request batch supports Qwen-Image, Flux, SD3.5, and LTX2.3. The implementation reuses/adapts the existing prompt-list batching path in pipelines, so we first enabled it for the models we care about most. |
Signed-off-by: jader <yjader@foxmail.com>
960df9c to
dcf55c0
Compare
|
I think this PR needs to polish for further abstraction, but let's wait for the next version |
…tch API PR #4079 replaced per-request objects with `DiffusionRequestBatch` (whose `sampling_params` is a read-only property) and `OmniDiffusionRequest` (which exposes a singular `prompt`, not `prompts`). The two-stage video pipelines were not migrated, so they crashed at worker dummy-run and the server never became ready: - ltx2/pipeline_ltx2.py, ltx2/pipeline_ltx2_image2video.py: `stage_2_req.sampling_params = req.sampling_params.clone()` hit "property 'sampling_params' of 'DiffusionRequestBatch' object has no setter". Now clone the underlying request(s) and override their sampling params, leaving the batch's read-only property alone. - helios/pipeline_helios.py: `prepare_encode` built a single OmniDiffusionRequest then read `req.prompts`, hitting "'OmniDiffusionRequest' object has no attribute 'prompts'". Now wrap the request in a DiffusionRequestBatch so the batch compatibility properties are available. Fixes the "Diffusion X2V - Other Function/Accuracy Test" failures (test_ltx2_expansion, test_video_streaming_output_similarity[helios_distilled]). Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com> Signed-off-by: tzhouam <tzhouam@connect.ust.hk>
Integrate upstream main 0899a1a (11 commits since 1b318d1) — orchestrator inter-stage/client output split (vllm-project#4527), diffusion request-level batching (vllm-project#4079), speech SSE default (vllm-project#4679), and assorted model/example migrations. Auto-merged with no conflicts. Signed-off-by: chickeyton <ngton2014@gmail.com> Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
…m-project#4079) Upstream vllm-project#4079 ("request-level batching for diffusion pipelines") made StagePool.submit_initial / submit_update reject diffusion list-prompt batch requests. The orchestrator's AR→DiT handoff (_forward_to_next_stage) only unwrapped a length-1 diffusion-prompt list on the streaming-update path, so the initial submit still handed a list to submit_initial and the orchestrator thread died on the first request reaching the DiT stage (EngineDeadError). Unwrap a length-1 list for both the initial and the update submit so each diffusion request is submitted independently. Multi-prompt batches (len > 1), which previously rode the removed list-prompt path, now raise a clear error since per-sub-request orchestration is not yet wired up. Validated e2e on HunyuanImage TI2I (C2: 2-rep TP2 EP-off, stage-based): 36/36 requests succeeded, 0 list-prompt errors (was 1/4 then 0/16 + crash). Signed-off-by: chickeyton <ngton2014@gmail.com> Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
Merging main pulled in new call sites that used the pre-rename output queue API, which this branch had already renamed output_async_queue -> output_sync_queue (put_nowait on the sync side). These were semantic conflicts: git auto-merged them without a textual conflict, so they landed referencing an attribute the constructor no longer sets. - orchestrator.py:1251 (from vllm-project#4079, diffusion request-level batching) and orchestrator.py:1419 (from vllm-project#4257, Aura non-async-chunk path) still called `await self.output_async_queue.put(...)`, which would raise AttributeError on those terminal-output error/edge paths. Convert to `self.output_sync_queue.put_nowait(...)`. - tests/engine/test_orchestrator_stage_input_bridge.py (new file from vllm-project#4257, marked core_model/cpu) constructed Orchestrator with the old `output_async_queue=` kwarg and failed against the new signature. Update to `output_sync_queue=output_q.sync_q`. Tested on vLLM 0.24.0: engine + orchestrator unit tests (41 passed) and a Qwen3-TTS-12Hz-0.6B streaming /v1/audio/speech smoke test. Signed-off-by: zeningc <zening.chen@yahoo.com>
…es (vllm-project#4079) Signed-off-by: jader <yjader@foxmail.com> Co-authored-by: Samit <285365963@qq.com> Signed-off-by: Wojciech Kutak <wkutak@nvidia.com>
…es (vllm-project#4079) Signed-off-by: jader <yjader@foxmail.com> Co-authored-by: Samit <285365963@qq.com>
Purpose
Implements Phase 2 of RFC #3550. This completes the diffusion request contract cleanup by making
OmniDiffusionRequestrepresent exactly one logical request with oneprompt, while runtime batching is represented by scheduler output and runner-side batch structures.In this PR:
Request contract
OmniDiffusionRequest.promptsis replaced by a singlepromptfield; request identity remains the required singlerequest_id.__post_init__rejects list inputs defensively.NewRequestData(request_id, req), carrying the already-initializedOmniDiffusionRequestdirectly so the executor/worker forward it instead of rebuilding it (which would re-trigger__post_init__and corrupt sentinel fields likeguidance_scale_provided). Fields likeprompt/sampling_params/kv_sender_infoare read offreqrather than duplicated on the payload.Scheduler batching
SamplingParamsKey, with conservative LoRA homogeneity constraints.num_outputs_per_promptis added to the key so requests producing differently-shaped outputs do not get co-batched.Request-level batch & execution path
DiffusionRequestBatchis introduced as the request-level batch abstraction above the step/tensor-levelInputBatch(StepInputBatchalias). It wrapslist[OmniDiffusionRequest]and exposesprompts/sampling_params/request_id/kv_sender_infocompatibility properties so pipelineforwardbodies stay close to upstream. Worker-sideDiffusionRequestState(single prompt) is retained for the stepwise path only.forwardis unified toforward(req: DiffusionRequestBatch) -> list[DiffusionOutput]. A class attributesupports_request_batchadvertises whether the pipeline supports one batchedforward.execute_request(oneexecute_modelRPC per scheduled request, the preserved upstream path) andexecute_batch(a singleexecute_model_batchRPC carrying the wholeDiffusionSchedulerOutput). The engine resolves request-batch capability at initialization from the configured pipeline class, including custom pipeline classes, and binds batch-capable pipelines directly toexecute_batchfor all request-mode cycles, including single-request cycles.execute_model_batchbuilds theDiffusionRequestBatchand runs per-request setup (KV transfer, generator/seed), one cache refresh, and one LoRA activation per batch. Non-batch pipelines stay on the per-request execution path.BatchRunnerOutputroutes results perrequest_id, including per-request success, error, and abort outputs.RunnerOutput.resultwrappers and nestedBatchRunnerOutput.runner_outputs[*].result, so request-mode batch outputs with large tensors continue to use the SHM transfer path instead of falling back to pickle IPC.request_batch_max_wait_ms. The engine can wait briefly before the firstschedule()of a wave so bursty compatible requests can accumulate, while the default0.0keeps the existing no-wait behavior. The scheduler exposes waiting/running queue counters for this logic, and the serve CLI / default diffusion stage config forward the option intoOmniDiffusionConfig.API boundary (breaking change)
AsyncOmni.generate(prompt=[...])rejects early with a clear error when a diffusion stage is present;StagePool.submit_initial()/submit_update()keep a defensive list check.Pipeline & test migration
DiffusionRequestBatchforwardcontract; non-batch models change only signature/return type, batch-capable models addsupports_request_batch = Trueplus per-request output splitting.output_formatter.pyalso follows the single-request contract: it readsrequest.prompt, emits oneOmniRequestOutputper request, and no longer keeps arequest.promptscompatibility branch.Compatibility/docs/release notes
AsyncOmni.generate(prompt=[...])packed-list behavior. Users should submit multiple independent requests to use automatic scheduler batching.max_num_seqs, and the optionalrequest_batch_max_wait_msadmission wait setting.Test Plan
Core diffusion request-mode regression
Focused coverage for request contract, scheduler batching, runner behavior,
engine routing, multiprocess dispatch, and stage process integration:
Touched backend and model coverage
Representative backend/model unit coverage for changed diffusion paths, excluding
hardware-dependent advanced model tests:
python -m pytest \ tests/diffusion/diffusion_backend/test_diffusers_backend.py \ tests/diffusion/models/dmd2/test_dmd2_request_sanitization.py \ tests/diffusion/models/dmd2/test_dmd2_scheduler.py \ tests/diffusion/models/flux2/test_flux2_klein_num_inference_steps.py \ tests/diffusion/models/ovis_image/test_ovis_image.py \ tests/diffusion/models/qwen_image/test_qwen_image_edit_plus.py \ tests/diffusion/models/qwen_image/test_qwen_image_max_sequence_length.py \ tests/diffusion/models/wan2_2/test_wan22_i2v_pipeline.py \ tests/diffusion/models/wan2_2/test_wan22_pipeline_diffuse.py \ tests/diffusion/models/wan2_2/test_wan22_vace_pipeline.py \ -m "not advanced_model" \ -qOutput formatter regression
Focused coverage for the single-prompt output formatting contract and the removal of the old multi-prompt compatibility path:
IPC coverage
Large-tensor SHM packing coverage for request-mode and step-pipeline outputs:
Request-batch capability routing
Focused coverage for the
supports_request_batchroute from engine dispatch toworker-side
RequestBatchexecution:Request-batch admission and CLI forwarding
Focused coverage for request-batch capability detection, admission wait,
scheduler queue counters, and the
request_batch_max_wait_msCLI/config path:Qwen-Image performance comparison
Compare Qwen-Image performance with the same setup:
512x512,20denoisingsteps,
FLASH_ATTN, single A100,1warmup run, and30measured runs. Thebaseline uses merge-base packed prompt-list; the candidate and StepScheduler
checks use two concurrent single-prompt requests with
max_num_seqs: 2.Residual old-contract check
Run a residual old-contract grep before submission:
rg "req\\.prompts|request\\.prompts|OmniDiffusionRequest\\(.*prompts|add_batch_request_async" \ vllm_omni testsExpected remaining
add_batch_request_asynchits should be limited to the abstract interface or explicit rejection stubs. Remainingreq.prompts/request.promptshits should be inDiffusionRequestBatchcompatibility properties, pre-processing helpers, comments, or rejection messages — not in single-request execution paths.Test Result
137 passed, 23 warnings in 66.91s177 passed, 2 deselected, 20 warnings in 3.37s10 passed, 16 warnings in 2.56s25 passed, 16 warnings in 3.87s72 passed, 20 warnings in 3.51s18 passed, 17 warnings in 0.84sadd_batch_request_async/req.prompts/request.promptshits are in compatibility or rejection paths; focused
prompts=check now only finds theStageDiffusionProc rejection test.
Static request-batch route checks:
supports_request_batch = True.step_executiondisabled, setsmax_num_seqs: 2.DiffusionEngineselectsexecute_batchfor request-batch-capable pipelines whenstep_executionis disabled.Qwen-Image request batching performance comparison:
512x512,20steps,FLASH_ATTN, single A100,1warmup run,30measured runsstep_executiondisabled. StepScheduler benchmark is intentionally not used for this PR result.b8f681747771f3c2Current request scheduler batch delta vs merge-base prompt-list: mean
+0.15%, p50+0.04%, p90+0.54%, amortized mean+0.15%. No significant performance regression is observed with the new request-level batching path.