From 0828b511fb7508f6e8930747fef6ebae76bfcbe1 Mon Sep 17 00:00:00 2001 From: EphraimBuddy Date: Mon, 19 Apr 2021 20:07:25 +0100 Subject: [PATCH 1/2] Parse Error 410 in kubernetes Watcher and return latest resource version Currently, when kubernetes watcher stream encounters Error 410('too old resource version'), it returns resource version '0' which is not the latest version. This 410 error contains the latest resource version in its message. This PR parses the message and return the latest resource version so that watcher can continue watching instead of returning '0' --- airflow/executors/kubernetes_executor.py | 32 ++++++++++++++++++--- tests/executors/test_kubernetes_executor.py | 16 +++++++++++ 2 files changed, 44 insertions(+), 4 deletions(-) diff --git a/airflow/executors/kubernetes_executor.py b/airflow/executors/kubernetes_executor.py index 707fe4fcdc4b8..daf5b8099b7ea 100644 --- a/airflow/executors/kubernetes_executor.py +++ b/airflow/executors/kubernetes_executor.py @@ -25,6 +25,7 @@ import functools import json import multiprocessing +import re import time from queue import Empty, Queue # pylint: disable=unused-import from typing import Any, Dict, List, Optional, Tuple @@ -173,11 +174,15 @@ def process_error(self, event: Any) -> str: self.log.error('Encountered Error response from k8s list namespaced pod stream => %s', event) raw_object = event['raw_object'] if raw_object['code'] == 410: - self.log.info( - 'Kubernetes resource version is too old, must reset to 0 => %s', (raw_object['message'],) + message = raw_object['message'] + latest_resource_version = self._parse_too_old_failure(message) + self.log.warning( + "Updated to new resource version: %s due to 'too old' error: %s", + latest_resource_version, + raw_object, ) - # Return resource version 0 - return '0' + + return latest_resource_version raise AirflowException( 'Kubernetes failure for %s with code %s and message: %s' % (raw_object['reason'], raw_object['code'], raw_object['message']) @@ -218,6 +223,25 @@ def process_status( resource_version, ) + def _parse_too_old_failure(self, message): + """ + Parse stream watcher 410 'too old resource version' error + to get the latest resource version from the message + + See https://github.com/kubernetes-client/python/issues/609 + """ + regex = r"too old resource version: .* \((.*)\)" + result = re.search(regex, message) + if result is None: + return None + match = result.group(1) + if match is None: + return None + try: + return match + except (ValueError, TypeError): + return None + class AirflowKubernetesScheduler(LoggingMixin): """Airflow Scheduler for Kubernetes""" diff --git a/tests/executors/test_kubernetes_executor.py b/tests/executors/test_kubernetes_executor.py index 0129b3af5cdd4..468ebb7134e63 100644 --- a/tests/executors/test_kubernetes_executor.py +++ b/tests/executors/test_kubernetes_executor.py @@ -507,3 +507,19 @@ def test_process_status_catchall(self): self._run() self.watcher.watcher_queue.put.assert_not_called() + + @mock.patch.object(KubernetesJobWatcher, '_parse_too_old_failure') + def test_process_error_event(self, mock_too_old_failure): + message = "too old resource version: 27272 (43334)" + mock_too_old_failure.return_value = '43334' + self.pod.status.phase = 'Pending' + self.pod.metadata.resource_version = '43334' + raw_object = {"code": 410, "message": message} + self.events.append({"type": "ERROR", "object": self.pod, "raw_object": raw_object}) + self._run() + mock_too_old_failure.assert_called_once_with(message) + + def test_parse_too_old_failure(self): + message = "too old resource version: 27272 (43334)" + latest_version = self.watcher._parse_too_old_failure(message) + assert latest_version == '43334' From cc8ae273736c4e447cc7b460c01f59de8a21c6ec Mon Sep 17 00:00:00 2001 From: EphraimBuddy Date: Mon, 19 Apr 2021 20:36:10 +0100 Subject: [PATCH 2/2] fixup! Parse Error 410 in kubernetes Watcher and return latest resource version --- airflow/executors/kubernetes_executor.py | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/airflow/executors/kubernetes_executor.py b/airflow/executors/kubernetes_executor.py index daf5b8099b7ea..c3475f28e3700 100644 --- a/airflow/executors/kubernetes_executor.py +++ b/airflow/executors/kubernetes_executor.py @@ -176,12 +176,14 @@ def process_error(self, event: Any) -> str: if raw_object['code'] == 410: message = raw_object['message'] latest_resource_version = self._parse_too_old_failure(message) + if latest_resource_version is None: + # Return resource version 0 + return '0' self.log.warning( "Updated to new resource version: %s due to 'too old' error: %s", latest_resource_version, raw_object, ) - return latest_resource_version raise AirflowException( 'Kubernetes failure for %s with code %s and message: %s'