Support controller in TransferQueue - #2
#2Support controller in TransferQueue#2
Conversation
There was a problem hiding this comment.
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) |
There was a problem hiding this comment.
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))
| 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)) |
| 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}") |
There was a problem hiding this comment.
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}')
| print(f"Data update error from storage: {error_msg}") | |
| logger.error(f"Data update error from storage: {error_msg}") |
| data_fields: list[str], | ||
| batch_size: int, | ||
| mode: str = "fetch", | ||
| global_step=0, |
There was a problem hiding this comment.
[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.
| 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) |
There was a problem hiding this comment.
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.
| 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 |
There was a problem hiding this comment.
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}'
| 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}" |
* Support controller in TransferQueue * Fix import * Fix comments --------- Co-authored-by: liuximeng <13073314+liuximeng18772102439@user.noreply.gitee.com>
…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>
What does this PR do?
Checklist Before Starting
[{modules}] {type}: {description}(This will be checked by the CI){modules}includefsdp,megatron,sglang,vllm,rollout,trainer,ci,training_utils,recipe,hardware,deployment,ray,worker,single_controller,misc,perf,model,algo,env,tool,ckpt,doc,data,like[megatron, fsdp, doc]{type}is infeat,fix,refactor,chore,test[BREAKING]to the beginning of the title.[BREAKING][fsdp, megatron] feat: dynamic batchingTest
API and Usage Example
# Add code snippet or script demonstrating how to use thisDesign & Code Changes
Checklist Before Submitting
Important
Please check all the following items before requesting a review, otherwise the reviewer might deprioritize this PR for review.
pre-commit install && pre-commit run --all-files --show-diff-on-failure --color=alwaysci-requestchannel in theverlSlack workspace. (If not accessible, please try the Feishu group (飞书群).)