Skip to content

Navigation Menu

Sign in
Appearance settings

Search code, repositories, users, issues, pull requests...

Provide feedback

We read every piece of feedback, and take your input very seriously.

Saved searches

Use saved searches to filter your results more quickly

Appearance settings

Support controller in TransferQueue - #2

#2
Merged
0oshowero0 merged 3 commits into
TransferQueue:mainTransferQueue/verl:mainfrom
LLLLxmmm:add_controllerLLLLxmmm/verl:add_controllerCopy head branch name to clipboard
Sep 24, 2025
Merged

Support controller in TransferQueue#2
0oshowero0 merged 3 commits into
TransferQueue:mainTransferQueue/verl:mainfrom
LLLLxmmm:add_controllerLLLLxmmm/verl:add_controllerCopy head branch name to clipboard

Conversation

@LLLLxmmm

Copy link
Copy Markdown

What does this PR do?

Add concise overview of what this PR aims to achieve or accomplish. Reference related GitHub issues and PRs that help with the review.

Checklist Before Starting

  • Search for similar PRs. Paste at least one query link here: ...
  • 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.

# 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.

@0oshowero0
0oshowero0 requested a review from Copilot September 23, 2025 13:54

Copilot AI left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Pull Request Overview

This PR introduces a new TransferQueueController class that enables controller functionality in the TransferQueue system. The controller manages data production and consumption status across distributed storage units, coordinating data flow between producers and consumers through ZMQ-based communication.

  • Adds centralized controller for managing data status and metadata across storage units
  • Implements ZMQ-based communication protocols for handshaking, request handling, and status updates
  • Provides data sampling and batch management capabilities with support for multiple consumption modes

Tip: Customize your code reviews with copilot-instructions.md. Create the file or learn how to get started.


TQ_CONTROLLER_GET_METADATA_TIMEOUT = int(os.environ.get("TQ_CONTROLLER_GET_METADATA_TIMEOUT", 300))
TQ_CONTROLLER_GET_METADATA_CHECK_INTERVAL = int(os.environ.get("TQ_CONTROLLER_GET_METADATA_CHECK_INTERVAL", 1))
TQ_INIT_FIELD_NUM = os.environ.get("TQ_INIT_FIELD_NUM", 10)

Copilot AI Sep 23, 2025

Copy link

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Environment variable TQ_INIT_FIELD_NUM is retrieved as a string but used as an integer. This will cause a TypeError when used in tensor operations. Convert to int: int(os.environ.get('TQ_INIT_FIELD_NUM', 10))

Suggested change
TQ_INIT_FIELD_NUM = os.environ.get("TQ_INIT_FIELD_NUM", 10)
TQ_INIT_FIELD_NUM = int(os.environ.get("TQ_INIT_FIELD_NUM", 10))

Copilot uses AI. Check for mistakes.
elif request_msg.request_type == ZMQRequestType.NOTIFY_DATA_UPDATE_ERROR:
# Handle data update errors
error_msg = request_msg.body.get("message", "Unknown error")
print(f"Data update error from storage: {error_msg}")

Copilot AI Sep 23, 2025

Copy link

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Using print() for error logging is inconsistent with the established logging pattern used throughout the file. Replace with logger.error(f'Data update error from storage: {error_msg}')

Suggested change
print(f"Data update error from storage: {error_msg}")
logger.error(f"Data update error from storage: {error_msg}")

Copilot uses AI. Check for mistakes.
data_fields: list[str],
batch_size: int,
mode: str = "fetch",
global_step=0,

Copilot AI Sep 23, 2025

Copy link

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[nitpick] The default value of 0 for global_step parameter could be problematic as it's not clear if this is a valid default for all use cases. Consider making this a required parameter or use None to indicate it must be explicitly provided.

Copilot uses AI. Check for mistakes.
current_fields = self.data_production_status.shape[1]
# Expand data status matrix if needed
if len(self.field_name_mapping) + needed_fields > current_fields:
add_fields = max(TQ_INIT_FIELD_NUM, needed_fields + 1)

Copilot AI Sep 23, 2025

Copy link

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

TQ_INIT_FIELD_NUM is a string from environment variable but used in max() comparison with integers. This will cause a TypeError. This is a secondary issue stemming from the type conversion problem in line 39.

Copilot uses AI. Check for mistakes.
if mode == "insert":
# TODO: Currently only supports putting entire GBS data, need to extend to support multiple puts to same
# step
assert batch_size == self.global_batch_size

Copilot AI Sep 23, 2025

Copy link

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The assert statement for batch_size validation should include a descriptive error message to help with debugging: assert batch_size == self.global_batch_size, f'batch_size {batch_size} must equal global_batch_size {self.global_batch_size}'

Suggested change
assert batch_size == self.global_batch_size
assert batch_size == self.global_batch_size, f"batch_size {batch_size} must equal global_batch_size {self.global_batch_size}"

Copilot uses AI. Check for mistakes.
@0oshowero0
0oshowero0 merged commit 501e1e2 into TransferQueue:main Sep 24, 2025
0oshowero0 pushed a commit that referenced this pull request Sep 28, 2025
* Support controller in TransferQueue

* Fix import

* Fix comments

---------

Co-authored-by: liuximeng <13073314+liuximeng18772102439@user.noreply.gitee.com>
ji-huazhong added a commit that referenced this pull request Nov 18, 2025
…oller

* Support storage unit in TransferQueue

* Fix importance error

* Support controller in TransferQueue (#2)

* Support controller in TransferQueue

* Fix import

* Fix comments

---------

Co-authored-by: liuximeng <13073314+liuximeng18772102439@user.noreply.gitee.com>

* expose TransferQueueClient (#3)

* Add copyright and license information

Added copyright and licensing information to the controller.py file.

* update client docstring (#5)

Signed-off-by: 0oshowero0 <o0shower0o@outlook.com>

* merge TransferQueue utils (#4)

* [fix] Fix n_sample related problems (#8)

* update client docstring

Signed-off-by: 0oshowero0 <o0shower0o@outlook.com>

* fix n_sample related problems

Signed-off-by: 0oshowero0 <o0shower0o@outlook.com>

---------

Signed-off-by: 0oshowero0 <o0shower0o@outlook.com>

* expose TransferQueue client/controller UT (#6)

* Add metadata.py and test_simple_storage_unit.py (#9)

* Add metadata.py and test_simple_storage_unit.py

* Add copyright and license information to test_simple_storage_unit.py

* Apply suggestion from @Copilot

Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com>

---------

Co-authored-by: Han Zhenyu 韩振宇 <o0shower0o@outlook.com>
Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com>

* Add reorder function to BatchMeta (#13)

Co-authored-by: liuximeng <13073314+liuximeng18772102439@user.noreply.gitee.com>

* [recipe, data] feat: TransferQueue - Support managing multiple data partitions for Train/Val/Test in controller (#45)

Co-authored-by: liuximeng <13073314+liuximeng18772102439@user.noreply.gitee.com>

* delete TQ source codes

Signed-off-by: 0oshowero0 <o0shower0o@outlook.com>

* update docs

Signed-off-by: 0oshowero0 <o0shower0o@outlook.com>

* update performance

Signed-off-by: 0oshowero0 <o0shower0o@outlook.com>

* fix

Signed-off-by: 0oshowero0 <o0shower0o@outlook.com>

---------

Signed-off-by: 0oshowero0 <o0shower0o@outlook.com>
Co-authored-by: FightingZhen <295632982@qq.com>
Co-authored-by: Han Zhenyu 韩振宇 <o0shower0o@outlook.com>
Co-authored-by: LLLLxmmm <130739718+LLLLxmmm@users.noreply.github.com>
Co-authored-by: liuximeng <13073314+liuximeng18772102439@user.noreply.gitee.com>
Co-authored-by: Han Zhenyu 韩振宇 <hanzy19@tsinghua.org.cn>
Co-authored-by: zhabuye <74179177+zhabuye@users.noreply.github.com>
Co-authored-by: Jianjun Zhong <87791082+jianjunzhong@users.noreply.github.com>
Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com>
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants

Morty Proxy This is a proxified and sanitized view of the page, visit original site.