Skip to content

Commit ecdaa8d

Browse files
authored
[recipe] fix: fix bug of tranfer queue runtime env (volcengine#3904)
### What does this PR do? as title ### Checklist Before Starting - [x] Search for similar PRs. Paste at least one query link here: ... - [x] Format the PR title as `[{modules}] {type}: {description}` (This will be checked by the CI) - `{modules}` include `fsdp`, `megatron`, `sglang`, `vllm`, `rollout`, `trainer`, `ci`, `training_utils`, `recipe`, `hardware`, `deployment`, `ray`, `worker`, `single_controller`, `misc`, `perf`, `model`, `algo`, `env`, `tool`, `ckpt`, `doc`, `data` - If this PR involves multiple modules, separate them with `,` like `[megatron, fsdp, doc]` - `{type}` is in `feat`, `fix`, `refactor`, `chore`, `test` - If this PR breaks any API (CLI arguments, config, function signature, etc.), add `[BREAKING]` to the beginning of the title. - Example: `[BREAKING][fsdp, megatron] feat: dynamic batching` ### Test > For changes that can not be tested by CI (e.g., algorithm implementation, new model support), validate by experiment(s) and show results like training curve plots, evaluation results, etc. ### API and Usage Example > Demonstrate how the API changes if any, and provide usage example(s) if possible. ```python # Add code snippet or script demonstrating how to use this ``` ### Design & Code Changes > Demonstrate the high-level design if this PR is complex, and list the specific changes. ### Checklist Before Submitting > [!IMPORTANT] > Please check all the following items before requesting a review, otherwise the reviewer might deprioritize this PR for review. - [x] Read the [Contribute Guide](https://github.com/volcengine/verl/blob/main/CONTRIBUTING.md). - [x] Apply [pre-commit checks](https://github.com/volcengine/verl/blob/main/CONTRIBUTING.md#code-linting-and-formatting): `pre-commit install && pre-commit run --all-files --show-diff-on-failure --color=always` - [ ] Add / Update [the documentation](https://github.com/volcengine/verl/tree/main/docs). - [ ] Add unit or end-to-end test(s) to [the CI workflow](https://github.com/volcengine/verl/tree/main/.github/workflows) to cover all the code. If not feasible, explain why: ... - [ ] Once your PR is ready for CI, send a message in [the `ci-request` channel](https://verl-project.slack.com/archives/C091TCESWB1) in [the `verl` Slack workspace](https://join.slack.com/t/verl-project/shared_invite/zt-3855yhg8g-CTkqXu~hKojPCmo7k_yXTQ). (If not accessible, please try [the Feishu group (飞书群)](https://applink.larkoffice.com/client/chat/chatter/add_by_link?link_token=772jd4f1-cd91-441e-a820-498c6614126a).)
1 parent ab4c04b commit ecdaa8d

File tree

2 files changed

+14
-4
lines changed

2 files changed

+14
-4
lines changed

recipe/transfer_queue/main_ppo.py

Lines changed: 7 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -67,10 +67,15 @@ def run_ppo(config, task_runner_class=None) -> None:
6767
default_runtime_env = get_ppo_ray_runtime_env()
6868
ray_init_kwargs = config.ray_kwargs.get("ray_init", {})
6969
runtime_env_kwargs = ray_init_kwargs.get("runtime_env", {})
70+
71+
if config.transfer_queue.enable:
72+
# Add runtime environment variables for transfer queue
73+
runtime_env_vars = runtime_env_kwargs.get("env_vars", {})
74+
runtime_env_vars["TRANSFER_QUEUE_ENABLE"] = "1"
75+
runtime_env_kwargs["env_vars"] = runtime_env_vars
76+
7077
runtime_env = OmegaConf.merge(default_runtime_env, runtime_env_kwargs)
7178
ray_init_kwargs = OmegaConf.create({**ray_init_kwargs, "runtime_env": runtime_env})
72-
if config.transfer_queue.enable:
73-
ray_init_kwargs["TRANSFER_QUEUE_ENABLE"] = "1"
7479
print(f"ray init kwargs: {ray_init_kwargs}")
7580
ray.init(**OmegaConf.to_container(ray_init_kwargs))
7681

verl/trainer/main_ppo.py

Lines changed: 7 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -61,10 +61,15 @@ def run_ppo(config, task_runner_class=None) -> None:
6161
default_runtime_env = get_ppo_ray_runtime_env()
6262
ray_init_kwargs = config.ray_kwargs.get("ray_init", {})
6363
runtime_env_kwargs = ray_init_kwargs.get("runtime_env", {})
64+
65+
if config.transfer_queue.enable:
66+
# Add runtime environment variables for transfer queue
67+
runtime_env_vars = runtime_env_kwargs.get("env_vars", {})
68+
runtime_env_vars["TRANSFER_QUEUE_ENABLE"] = "1"
69+
runtime_env_kwargs["env_vars"] = runtime_env_vars
70+
6471
runtime_env = OmegaConf.merge(default_runtime_env, runtime_env_kwargs)
6572
ray_init_kwargs = OmegaConf.create({**ray_init_kwargs, "runtime_env": runtime_env})
66-
if config.transfer_queue.enable:
67-
ray_init_kwargs["TRANSFER_QUEUE_ENABLE"] = "1"
6873
print(f"ray init kwargs: {ray_init_kwargs}")
6974
ray.init(**OmegaConf.to_container(ray_init_kwargs))
7075

0 commit comments

Comments
 (0)