From 9f252d246e734a0ecb3870b551175e6d0c19c7b4 Mon Sep 17 00:00:00 2001 From: Alexander Kukushkin Date: Thu, 11 Feb 2021 15:55:05 +0100 Subject: [PATCH] Improve handling of concurrent update error (#1796) The old strategy was waiting for 1 second and hoping that we will get an update event from the WATCH connection. Unfortunately, it didn't work well in practice. Instead, we will get the current value from the API by performing an explicit read request. Close https://github.com/zalando/patroni/issues/1767 --- patroni/dcs/kubernetes.py | 19 ++++++++++++++----- tests/test_kubernetes.py | 28 ++++++++++++++++++---------- 2 files changed, 32 insertions(+), 15 deletions(-) diff --git a/patroni/dcs/kubernetes.py b/patroni/dcs/kubernetes.py index a0a08978..abd98a54 100644 --- a/patroni/dcs/kubernetes.py +++ b/patroni/dcs/kubernetes.py @@ -905,13 +905,23 @@ class Kubernetes(AbstractDCS): except (RetryFailedError, K8sException): return False - deadline = retry.stoptime - time.time() - if deadline < 2: + retry.deadline = retry.stoptime - time.time() + if retry.deadline < 1: return False - retry.sleep_func(1) # Give a chance for ObjectCache to receive the latest version + # Try to get the latest version directly from K8s API instead of relying on async cache + try: + kind = retry(self._api.read_namespaced_kind, self.leader_path, self._namespace) + except Exception as e: + logger.error('Failed to get the leader object "%s": %r', self.leader_path, e) + return False + + self._kinds.set(self.leader_path, kind) + + retry.deadline = retry.stoptime - time.time() + if retry.deadline < 0.5: + return False - kind = self._kinds.get(self.leader_path) kind_annotations = kind and kind.metadata.annotations or {} kind_resource_version = kind and kind.metadata.resource_version @@ -919,7 +929,6 @@ class Kubernetes(AbstractDCS): if kind and (kind_annotations.get(self._LEADER) != self._name or kind_resource_version == resource_version): return False - retry.deadline = deadline - 1 # Update deadline and retry return self.patch_or_create(self.leader_path, annotations, kind_resource_version, ips=ips, retry=_retry) def update_leader(self, last_operation, access_is_restricted=False): diff --git a/tests/test_kubernetes.py b/tests/test_kubernetes.py index cb9462e6..aa4493a9 100644 --- a/tests/test_kubernetes.py +++ b/tests/test_kubernetes.py @@ -26,7 +26,7 @@ def mock_list_namespaced_config_map(*args, **kwargs): return k8s_client.V1ConfigMapList(metadata=metadata, items=items, kind='ConfigMapList') -def mock_list_namespaced_endpoints(*args, **kwargs): +def mock_read_namespaced_endpoints(*args, **kwargs): target_ref = k8s_client.V1ObjectReference(kind='Pod', resource_version='10', name='p-0', namespace='default', uid='964dfeae-e79b-4476-8a5a-1920b5c2a69d') address0 = k8s_client.V1EndpointAddress(ip='10.0.0.0', target_ref=target_ref) @@ -35,9 +35,12 @@ def mock_list_namespaced_endpoints(*args, **kwargs): subset = k8s_client.V1EndpointSubset(addresses=[address1, address0], ports=[port]) metadata = k8s_client.V1ObjectMeta(resource_version='1', labels={'f': 'b'}, name='test', annotations={'optime': '1234', 'leader': 'p-0', 'ttl': '30s'}) - endpoint = k8s_client.V1Endpoints(subsets=[subset], metadata=metadata) - metadata = k8s_client.V1ObjectMeta(resource_version='1') - return k8s_client.V1EndpointsList(metadata=metadata, items=[endpoint], kind='V1EndpointsList') + return k8s_client.V1Endpoints(subsets=[subset], metadata=metadata) + + +def mock_list_namespaced_endpoints(*args, **kwargs): + return k8s_client.V1EndpointsList(metadata=k8s_client.V1ObjectMeta(resource_version='1'), + items=[mock_read_namespaced_endpoints()], kind='V1EndpointsList') def mock_list_namespaced_pod(*args, **kwargs): @@ -261,20 +264,25 @@ class TestKubernetesEndpoints(BaseTestKubernetes): def test_update_leader_with_restricted_access(self): self.assertIsNotNone(self.k.update_leader('123', True)) + @patch.object(k8s_client.CoreV1Api, 'read_namespaced_endpoints', create=True) @patch.object(k8s_client.CoreV1Api, 'patch_namespaced_endpoints', create=True) - def test__update_leader_with_retry(self, mock_patch): + def test__update_leader_with_retry(self, mock_patch, mock_read): + mock_read.return_value = mock_read_namespaced_endpoints() mock_patch.side_effect = k8s_client.rest.ApiException(502, '') self.assertFalse(self.k.update_leader('123')) mock_patch.side_effect = RetryFailedError('') self.assertFalse(self.k.update_leader('123')) mock_patch.side_effect = k8s_client.rest.ApiException(409, '') - with patch('time.time', Mock(side_effect=[0, 100, 200])): + with patch('time.time', Mock(side_effect=[0, 100, 200, 0, 0, 0, 0, 100, 200])): self.assertFalse(self.k.update_leader('123')) - with patch('time.sleep', Mock()): self.assertFalse(self.k.update_leader('123')) - mock_patch.side_effect = [k8s_client.rest.ApiException(409, ''), mock_namespaced_kind()] - self.k._kinds._object_cache['test'].metadata.resource_version = '2' - self.assertIsNotNone(self.k._update_leader_with_retry({}, '1', [])) + self.assertFalse(self.k.update_leader('123')) + mock_patch.side_effect = [k8s_client.rest.ApiException(409, ''), mock_namespaced_kind()] + mock_read.return_value.metadata.resource_version = '2' + self.assertIsNotNone(self.k._update_leader_with_retry({}, '1', [])) + mock_patch.side_effect = k8s_client.rest.ApiException(409, '') + mock_read.side_effect = Exception + self.assertFalse(self.k.update_leader('123')) @patch.object(k8s_client.CoreV1Api, 'create_namespaced_endpoints', Mock(side_effect=[k8s_client.rest.ApiException(500, ''),