From cd871244e04e77b6e3cd81ffc96a275f3fbabfca Mon Sep 17 00:00:00 2001 From: Milly Date: Wed, 25 Apr 2018 02:35:54 +0900 Subject: [PATCH] `TransportError` not raised in `helpers.streaming_bulk()` with `max_retries`. (#775) * Add test to raise `TransportError` in `helpers.streaming_bulk()` with `max_retries`. * Re-raise `TransportError` if the last retry fails. --- elasticsearch/helpers/__init__.py | 2 +- test_elasticsearch/test_server/test_helpers.py | 17 +++++++++++++++++ 2 files changed, 18 insertions(+), 1 deletion(-) diff --git a/elasticsearch/helpers/__init__.py b/elasticsearch/helpers/__init__.py index ad741de5..eb932999 100644 --- a/elasticsearch/helpers/__init__.py +++ b/elasticsearch/helpers/__init__.py @@ -210,7 +210,7 @@ def streaming_bulk(client, actions, chunk_size=500, max_chunk_bytes=100 * 1024 * except TransportError as e: # suppress 429 errors since we will retry them - if not max_retries or e.status_code != 429: + if attempt == max_retries or e.status_code != 429: raise else: if not to_retry: diff --git a/test_elasticsearch/test_server/test_helpers.py b/test_elasticsearch/test_server/test_helpers.py index 2e27a681..2190a05f 100644 --- a/test_elasticsearch/test_server/test_helpers.py +++ b/test_elasticsearch/test_server/test_helpers.py @@ -136,6 +136,23 @@ class TestStreamingBulk(ElasticsearchTestCase): self.assertEquals(2, res['hits']['total']) self.assertEquals(4, failing_client._called) + def test_transport_error_is_raised_with_max_retries(self): + failing_client = FailingBulkClient(self.client, fail_at=(1, 2, 3, 4, ), + fail_with=TransportError(429, 'Rejected!', {})) + + def streaming_bulk(): + results = list(helpers.streaming_bulk( + failing_client, + [{"a": 42}, {"a": 39}], + raise_on_exception=True, + max_retries=3, + initial_backoff=0 + )) + return results + + self.assertRaises(TransportError, streaming_bulk) + self.assertEquals(4, failing_client._called) + class TestBulk(ElasticsearchTestCase): def test_bulk_works_with_single_item(self):