| 4 | from .k8s import K8s |
| 5 | |
| 6 | class MasterExecutor(TaskExecutor): |
| 7 | |
| 8 | def load(self): |
| 9 | self.lock = Lock() |
| 10 | self.k8s = K8s() |
| 11 | # Initialize capacities |
| 12 | self.current_capacity = {"request": 0, "browser": 0, "task": 0} |
| 13 | self.max_capacity = self.calculate_max_capacity() |
| 14 | |
| 15 | def get_max_running_count(self, scraper_type): |
| 16 | return self.max_capacity[scraper_type] |
| 17 | |
| 18 | def get_current_running_count(self, scraper_type): |
| 19 | return self.current_capacity[scraper_type] |
| 20 | |
| 21 | def run_task_and_update_state(self, scraper_type, task_json): |
| 22 | node = self.k8s.get_min_capacity_node(scraper_type) |
| 23 | |
| 24 | # Send task to the selected node |
| 25 | self.k8s.run_task_on_node(task_json, node['node_name']) |
| 26 | |
| 27 | # Update capacities |
| 28 | self.increment_master_capacity(node, scraper_type) |
| 29 | |
| 30 | def increment_master_capacity(self, node , scraper_type): |
| 31 | with self.lock: |
| 32 | node['current_capacity'][scraper_type] += 1 |
| 33 | self.current_capacity[scraper_type] += 1 |
| 34 | |
| 35 | def decrement_master_capacity(self, node , scraper_type): |
| 36 | with self.lock: |
| 37 | node['current_capacity'][scraper_type] -= 1 |
| 38 | self.current_capacity[scraper_type] -= 1 |
| 39 | |
| 40 | def calculate_max_capacity(self): |
| 41 | limit = Server.get_rate_limit() |
| 42 | replicas = len(self.k8s.nodes) |
| 43 | return {"request": limit["request"] * replicas,"task": limit["task"] * replicas, "browser": limit["browser"] * replicas} |
| 44 | |
| 45 | def on_success(self, task_id, task_type, task_result, node_name,scraper_name, data): |
| 46 | node = self.k8s.get_node(node_name) |
| 47 | self.decrement_master_capacity(node, task_type) |
| 48 | |
| 49 | # Further processing of task_result can be done here |
| 50 | self.mark_task_as_success(task_id, task_result, Server.cache,scraper_name, data) |
| 51 | self.update_parent_task(task_id, task_result) |
| 52 | |
| 53 | def on_failure(self, task_id, task_type, task_result, node_name): |
| 54 | node = self.k8s.get_node(node_name) |
| 55 | self.decrement_master_capacity(node, task_type) |
| 56 | |
| 57 | # Further processing of task_result can be done here |
| 58 | self.mark_task_as_failure(task_id, task_result) |
| 59 | self.update_parent_task(task_id, []) |