(self)
| 95 | return self._workdir |
| 96 | |
| 97 | def start(self): |
| 98 | self.state_start() |
| 99 | logging.info("Starting Dask cluster") |
| 100 | |
| 101 | server_port = find_free_port() |
| 102 | |
| 103 | # Make sure to kill all previously running Dask instances |
| 104 | subprocess.run(["killall", "dask"], check=False) |
| 105 | |
| 106 | self._start_scheduler(server_port) |
| 107 | |
| 108 | server_address = f"{HOSTNAME}:{server_port}" |
| 109 | self._start_workers(self.worker_assignment, server_address) |
| 110 | |
| 111 | self.client = Client(server_address) |
| 112 | self.client.wait_for_workers(n_workers=self.worker_count) |
| 113 | time.sleep(1) |
| 114 | assert len(self.client.scheduler_info()["workers"]) == self.worker_count |
| 115 | |
| 116 | self.cluster.start_monitoring(self.info.cluster_info.node_list.resolve(), observe_processes=True) |
| 117 | self.cluster.commit() |
| 118 | logging.info("Dask cluster started") |
| 119 | |
| 120 | def _start_scheduler(self, server_port: int): |
| 121 | scheduler = StartProcessArgs( |
nothing calls this directly
no test coverage detected