| 29 | delay = min(delay * backoff_factor, max_delay) + random.uniform(0, 0.1 * delay) # Adding jitter |
| 30 | |
| 31 | class WorkerExecutor(): |
| 32 | def load(self): |
| 33 | pass |
| 34 | |
| 35 | def start(self): |
| 36 | pass |
| 37 | |
| 38 | def run_worker_task(self, task, node_name): |
| 39 | Thread(target=self.perform_worker_task, args=(task, node_name,), daemon=True).start() |
| 40 | |
| 41 | def perform_worker_task(self, task, node_name): |
| 42 | scraper_type = task["scraper_type"] |
| 43 | task_id = task["id"] |
| 44 | scraper_name = task["scraper_name"] |
| 45 | metadata = {"metadata": task["metadata"]} if task["metadata"] != {} else {} |
| 46 | task_data = task["data"] |
| 47 | |
| 48 | fn = Server.get_scraping_function(scraper_name) |
| 49 | |
| 50 | try: |
| 51 | result = fn( |
| 52 | task_data, |
| 53 | **metadata, |
| 54 | parallel=None, |
| 55 | cache=False, |
| 56 | beep=False, |
| 57 | run_async=False, |
| 58 | async_queue=False, |
| 59 | raise_exception=True, |
| 60 | close_on_crash=True, |
| 61 | output=None, |
| 62 | create_error_logs=False |
| 63 | ) |
| 64 | result = clean_data(result) |
| 65 | make_request_with_retry(lambda: requests.post('http://master-srv:8000/k8s/success', json={"task_id": task_id, "task_type" : scraper_type, "data": task_data , "scraper_name": scraper_name , "task_result": result , "node_name": node_name})) |
| 66 | except Exception: |
| 67 | exception_log = traceback.format_exc() |
| 68 | traceback.print_exc() |
| 69 | make_request_with_retry(lambda: requests.post('http://master-srv:8000/k8s/fail', json={"task_id": task_id, "task_type" : scraper_type, "task_result": exception_log , "node_name": node_name})) |