Initializes the engine and starts its sub-services. If `api_server_pid` is defined, will launch a thread to keep getting request from zmq_server.
(
self, ipc_signal_suffix, local_data_parallel_id, request_queues_for_dp_ipc=None, result_queues_for_dp_ipc=None
)
| 97 | self._finalizer = weakref.finalize(self, self._exit_sub_services) |
| 98 | |
| 99 | def start( |
| 100 | self, ipc_signal_suffix, local_data_parallel_id, request_queues_for_dp_ipc=None, result_queues_for_dp_ipc=None |
| 101 | ): |
| 102 | """ |
| 103 | Initializes the engine and starts its sub-services. |
| 104 | If `api_server_pid` is defined, will launch a thread |
| 105 | to keep getting request from zmq_server. |
| 106 | """ |
| 107 | # assert not self.is_started, "The engine is already started." |
| 108 | |
| 109 | start_time = time.time() |
| 110 | self.engine.start() |
| 111 | if envs.FD_ENABLE_RETURN_TEXT: |
| 112 | self.engine.create_data_processor() |
| 113 | if self.cfg.scheduler_config.name == "dp": |
| 114 | self.cfg.init_cache_info() |
| 115 | assert (request_queues_for_dp_ipc is not None) and (result_queues_for_dp_ipc is not None) |
| 116 | self.engine.scheduler.start(local_data_parallel_id, request_queues_for_dp_ipc, result_queues_for_dp_ipc) |
| 117 | |
| 118 | if ipc_signal_suffix is not None: |
| 119 | self.api_server_pid = ipc_signal_suffix |
| 120 | self.engine.start_zmq_service(ipc_signal_suffix) |
| 121 | else: |
| 122 | ipc_signal_suffix = self.cfg.parallel_config.engine_worker_queue_port[0] |
| 123 | self.engine.start_zmq_service(self.cfg.parallel_config.engine_worker_queue_port[local_data_parallel_id]) |
| 124 | |
| 125 | self.llm_logger.info(f"start expert service {local_data_parallel_id}") |
| 126 | |
| 127 | if self.cfg.scheduler_config.name == "splitwise": |
| 128 | self.cfg.init_cache_info() |
| 129 | role = self.cfg.scheduler_config.splitwise_role |
| 130 | host_ip = self.cfg.host_ip |
| 131 | self.engine.scheduler.start(role, host_ip, self.cfg.register_info) |
| 132 | |
| 133 | if self.cfg.scheduler_config.splitwise_role != "mixed": |
| 134 | self.splitwise_receive_thread = threading.Thread( |
| 135 | target=self.engine.split_connector.start_receiver, args=() |
| 136 | ) |
| 137 | self.splitwise_receive_thread.daemon = True |
| 138 | self.splitwise_receive_thread.start() |
| 139 | self.cfg.print() |
| 140 | local_rank = local_data_parallel_id % self.cfg.worker_num_per_node |
| 141 | |
| 142 | if not envs.FD_ENABLE_MULTI_API_SERVER: |
| 143 | if self.cfg.parallel_config.data_parallel_size > 1: |
| 144 | launched_expert_service_signal_data = np.zeros( |
| 145 | shape=[self.cfg.parallel_config.data_parallel_size // self.cfg.nnode], dtype=np.int32 |
| 146 | ) |
| 147 | self.launched_expert_service_signal = IPCSignal( |
| 148 | name="launched_expert_service_signal", |
| 149 | array=launched_expert_service_signal_data, |
| 150 | dtype=np.int32, |
| 151 | suffix=ipc_signal_suffix, |
| 152 | create=False, |
| 153 | ) |
| 154 | self.launched_expert_service_signal.value[local_rank] = 1 |
| 155 | |
| 156 | if self.do_profile: |