启动所有线程并运行主循环
(self)
| 600 | f"create time: {oldest_task.status.created_time}") |
| 601 | |
| 602 | def start(self): |
| 603 | """启动所有线程并运行主循环""" |
| 604 | # 启动所有线程 |
| 605 | threads = [ |
| 606 | threading.Thread(target=self.reprocess, daemon=True), |
| 607 | threading.Thread(target=self.provider, daemon=True), |
| 608 | threading.Thread(target=self.reporter, daemon=True), |
| 609 | threading.Thread(target=self.serve_stealing, daemon=True), |
| 610 | threading.Thread(target=self.work_stealing, daemon=True), |
| 611 | threading.Thread(target=self.sync_handler, daemon=True), |
| 612 | threading.Thread(target=self.monitor, daemon=True) |
| 613 | ] |
| 614 | |
| 615 | for thread in threads: |
| 616 | thread.start() |
| 617 | |
| 618 | # 主循环 |
| 619 | while True: |
| 620 | tac: TaskAndContent = self.balance_collect.recv_pyobj() |
| 621 | self.cached_tasks[tac.status.completion_id] = tac |
| 622 | self.balance_collect.send_string("Received") |
| 623 | |
| 624 | # 发送任务状态到全局 |
| 625 | self.result_sender.send_pyobj(tac.status) |
| 626 | self.result_sender.recv_string() |
| 627 | |
| 628 | if len(self.cached_tasks) >= self.max_cache_size: |
| 629 | logger.warning(f"Too many cached tasks. [ Local GID: {self.local_gid} | " |
| 630 | f"Cached: {len(self.cached_tasks)}, Reprocessing: {self.valid_tasks.qsize()} | " |
| 631 | f"Queued: {self.ready_queue.qsize()} ]") |
| 632 | time.sleep(3) # 缓存过多 |
no test coverage detected