From 75dafcb558e67848bfd41d57b223b3a500001ca6 Mon Sep 17 00:00:00 2001 From: Frank Haldi <177933469+fhaldi@users.noreply.github.com> Date: Thu, 1 Oct 2026 16:29:08 +0200 Subject: [PATCH 1/6] Allow to continue on partial shard failures --- elastalert/elastalert.py | 20 ++++++++++++-------- 1 file changed, 12 insertions(+), 8 deletions(-) diff --git a/elastalert/elastalert.py b/elastalert/elastalert.py index b9a6be93..1f66f4c2 100755 --- a/elastalert/elastalert.py +++ b/elastalert/elastalert.py @@ -404,14 +404,18 @@ def get_hits(self, rule, starttime, endtime, index, scroll=False): self.thread_data.total_hits = int(res['hits']['total']['value']) - if len(res.get('_shards', {}).get('failures', [])) > 0: - try: - errs = [e['reason']['reason'] for e in res['_shards']['failures'] if 'Failed to parse' in e['reason']['reason']] - if len(errs): - raise ElasticsearchException(errs) - except (TypeError, KeyError): - # Different versions of ES have this formatted in different ways. Fallback to str-ing the whole thing - raise ElasticsearchException(str(res['_shards']['failures'])) + # Handles shard failures + shard_failures = res.get('_shards', {}).get('failures', []) + if shard_failures: + reasons = [str((f.get('reason') or {}).get('reason')) for f in shard_failures] + # Trigger Exception only if Elasticsearch query failed to parse + parse_errs = [r for r in reasons if 'Failed to parse' in r] + if parse_errs: + raise ElasticsearchException(parse_errs) + # If query still worked but returns some errors, continue... + elastalert_logger.warning( + 'Rule %s: %d shard(s) failed, continuing with partial results: %s', + rule['name'], len(shard_failures), shard_failures) elastalert_logger.debug(str(res)) except ElasticsearchException as e: From 99acb4c55e39110d0a91a10f85dac3f83da1a3fb Mon Sep 17 00:00:00 2001 From: Frank Haldi <177933469+fhaldi@users.noreply.github.com> Date: Thu, 1 Oct 2026 16:55:43 +0200 Subject: [PATCH 2/6] Allow to continue on partial shard failures --- elastalert/elastalert.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/elastalert/elastalert.py b/elastalert/elastalert.py index 1f66f4c2..471fc55d 100755 --- a/elastalert/elastalert.py +++ b/elastalert/elastalert.py @@ -406,7 +406,7 @@ def get_hits(self, rule, starttime, endtime, index, scroll=False): # Handles shard failures shard_failures = res.get('_shards', {}).get('failures', []) - if shard_failures: + if len(shard_failures) > 0: reasons = [str((f.get('reason') or {}).get('reason')) for f in shard_failures] # Trigger Exception only if Elasticsearch query failed to parse parse_errs = [r for r in reasons if 'Failed to parse' in r] From 9da01836d11c8fcd9b2f08a1cb9c45e1059a2d87 Mon Sep 17 00:00:00 2001 From: Frank Haldi <177933469+fhaldi@users.noreply.github.com> Date: Fri, 2 Oct 2026 12:23:42 +0200 Subject: [PATCH 3/6] Add doc, opt-in flag --- CHANGELOG.md | 1 + docs/source/configuration.rst | 3 +++ docs/source/ruletypes.rst | 7 ++++++- elastalert/config.py | 1 + elastalert/elastalert.py | 35 ++++++++++++++++++++++++----------- elastalert/test_rule.py | 1 + 6 files changed, 36 insertions(+), 12 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index a192f58d..01c5fd8d 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -5,6 +5,7 @@ ## New features - [PagerDuty] Add NTLM authentication support for HTTPS proxies - [#1772](https://github.com/jertel/elastalert2/pull/1772) - @sauravnz +- [Elasticsearch] Add opt-in flag to allow queries on degraded Elasticsearch indices - [#1777](https://github.com/jertel/elastalert2/pull/1777) - @fhaldi ## Other changes - [Helm] Omit enabled flag in probe definitions - [#1775](https://github.com/jertel/elastalert2/pull/1775) - @jim-barber-he diff --git a/docs/source/configuration.rst b/docs/source/configuration.rst index e822979a..f8a6e7e7 100644 --- a/docs/source/configuration.rst +++ b/docs/source/configuration.rst @@ -91,6 +91,9 @@ rule will no longer be run until either ElastAlert 2 restarts or the rule file h ``show_disabled_rules``: If true, ElastAlert 2 show the disable rules' list when finishes the execution. This defaults to True. +``allow_queries_on_degraded_indices``: If true, Elastalert 2 will still query Elasticsearch when some shards of the +data stream are unavailable. The default is ``false``. + ``notify_alert``: List of alerters to execute upon encountering a system error. System errors occur when an unexpected exception is thrown during rule processing. For additional notifications, such as when ElastAlert 2 background tests encounter problems, or when connectivity to the data storage system is lost, enable ``notify_all_errors``. See the :ref:`Alerts` section for the list of available alerters and their parameters. diff --git a/docs/source/ruletypes.rst b/docs/source/ruletypes.rst index 9da5edce..51735dfb 100644 --- a/docs/source/ruletypes.rst +++ b/docs/source/ruletypes.rst @@ -379,7 +379,6 @@ ca_certs ``ca_certs``: Path to a CA cert bundle to use to verify SSL connections (Optional, string, no default) - disable_rules_on_error ^^^^^^^^^^^^^^^^^^^^^^ @@ -387,6 +386,12 @@ disable_rules_on_error will upload a traceback message to ``elastalert_metadata`` and if ``notify_email`` is set, send an email notification. The rule will no longer be run until either ElastAlert 2 restarts or the rule file has been modified. This defaults to ``True``. +allow_queries_on_degraded_indices +^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ + +``allow_queries_on_degraded_indices``: If true, Elastalert 2 will still query Elasticsearch when some shards of the +data stream are unavailable. The default is ``false``. + es_conn_timeout ^^^^^^^^^^^^^^^ diff --git a/elastalert/config.py b/elastalert/config.py index 3d6b458f..b0387d56 100644 --- a/elastalert/config.py +++ b/elastalert/config.py @@ -76,6 +76,7 @@ def load_conf(args, defaults=None, overrides=None): conf.setdefault('scroll_keepalive', '30s') conf.setdefault('max_scrolling_count', 990) # Avoid stack overflow in run_query, note that 1000 is Python's stack limit conf.setdefault('disable_rules_on_error', True) + conf.setdefault('allow_queries_on_degraded_indices', False) conf.setdefault('scan_subdirectories', True) conf.setdefault('rules_loader', 'file') conf.setdefault('custom_pretty_ts_format', None) diff --git a/elastalert/elastalert.py b/elastalert/elastalert.py index 471fc55d..162ee985 100755 --- a/elastalert/elastalert.py +++ b/elastalert/elastalert.py @@ -145,6 +145,7 @@ def __init__(self, args): self.alert_time_limit = self.conf['alert_time_limit'] self.old_query_limit = self.conf['old_query_limit'] self.disable_rules_on_error = self.conf['disable_rules_on_error'] + self.allow_queries_on_degraded_indices = self.conf['allow_queries_on_degraded_indices'] self.notify_email = self.conf.get('notify_email', []) self.notify_all_errors = self.conf.get('notify_all_errors', False) self.notify_alert = self.conf.get('notify_alert', []) @@ -405,17 +406,29 @@ def get_hits(self, rule, starttime, endtime, index, scroll=False): self.thread_data.total_hits = int(res['hits']['total']['value']) # Handles shard failures - shard_failures = res.get('_shards', {}).get('failures', []) - if len(shard_failures) > 0: - reasons = [str((f.get('reason') or {}).get('reason')) for f in shard_failures] - # Trigger Exception only if Elasticsearch query failed to parse - parse_errs = [r for r in reasons if 'Failed to parse' in r] - if parse_errs: - raise ElasticsearchException(parse_errs) - # If query still worked but returns some errors, continue... - elastalert_logger.warning( - 'Rule %s: %d shard(s) failed, continuing with partial results: %s', - rule['name'], len(shard_failures), shard_failures) + if self.allow_queries_on_degraded_indices: + # Allow queries on degraded data streams + shard_failures = res.get('_shards', {}).get('failures', []) + if len(shard_failures) > 0: + reasons = [str((f.get('reason') or {}).get('reason')) for f in shard_failures] + # Trigger Exception only if Elasticsearch query failed to parse + parse_errs = [r for r in reasons if 'Failed to parse' in r] + if parse_errs: + raise ElasticsearchException(parse_errs) + # If query still worked but returns some errors, continue... + elastalert_logger.warning( + 'Rule %s: %d shard(s) failed, continuing with partial results: %s', + rule['name'], len(shard_failures), shard_failures) + else: + # Stop rule if any error with at least one shard of a data stream + if len(res.get('_shards', {}).get('failures', [])) > 0: + try: + errs = [e['reason']['reason'] for e in res['_shards']['failures'] if 'Failed to parse' in e['reason']['reason']] + if len(errs): + raise ElasticsearchException(errs) + except (TypeError, KeyError): + # Different versions of ES have this formatted in different ways. Fallback to str-ing the whole thing + raise ElasticsearchException(str(res['_shards']['failures'])) elastalert_logger.debug(str(res)) except ElasticsearchException as e: diff --git a/elastalert/test_rule.py b/elastalert/test_rule.py index a3b99f0a..3908d662 100644 --- a/elastalert/test_rule.py +++ b/elastalert/test_rule.py @@ -449,6 +449,7 @@ def run_rule_test(self): 'old_query_limit': {'weeks': 1}, 'run_every': {'minutes': 5}, 'disable_rules_on_error': False, + 'allow_queries_on_degraded_indices': False, 'buffer_time': {'minutes': 45}, 'scroll_keepalive': '30s' } From dfabd90078e2ae5e19711c76fa03129972128c52 Mon Sep 17 00:00:00 2001 From: Frank Haldi <177933469+fhaldi@users.noreply.github.com> Date: Fri, 2 Oct 2026 14:53:03 +0200 Subject: [PATCH 4/6] Add tests + clean warning if some shards failures --- elastalert/elastalert.py | 5 +++-- tests/base_test.py | 20 ++++++++++++++++++++ tests/config_test.py | 1 + tests/conftest.py | 1 + 4 files changed, 25 insertions(+), 2 deletions(-) diff --git a/elastalert/elastalert.py b/elastalert/elastalert.py index 162ee985..fcf8340c 100755 --- a/elastalert/elastalert.py +++ b/elastalert/elastalert.py @@ -410,15 +410,16 @@ def get_hits(self, rule, starttime, endtime, index, scroll=False): # Allow queries on degraded data streams shard_failures = res.get('_shards', {}).get('failures', []) if len(shard_failures) > 0: - reasons = [str((f.get('reason') or {}).get('reason')) for f in shard_failures] + reasons = [str(f) for f in shard_failures] # Trigger Exception only if Elasticsearch query failed to parse parse_errs = [r for r in reasons if 'Failed to parse' in r] if parse_errs: raise ElasticsearchException(parse_errs) # If query still worked but returns some errors, continue... + failed_indices = sorted({str(f.get('index')) for f in shard_failures if isinstance(f, dict)}) elastalert_logger.warning( 'Rule %s: %d shard(s) failed, continuing with partial results: %s', - rule['name'], len(shard_failures), shard_failures) + rule['name'], len(shard_failures), ', '.join(failed_indices)) else: # Stop rule if any error with at least one shard of a data stream if len(res.get('_shards', {}).get('failures', [])) > 0: diff --git a/tests/base_test.py b/tests/base_test.py index 6a345b61..79abdbd7 100644 --- a/tests/base_test.py +++ b/tests/base_test.py @@ -161,6 +161,26 @@ def test_no_hits(ea): assert ea.rules[0]['type'].add_data.call_count == 0 +@pytest.mark.parametrize('allow_degraded, failures, expected_success', [ + # All shards working: the option changes nothing + (False, [], True), + (True, [], True), + # Default: an unavailable shard stops the run + (False, [{'index': '.ds-index-test-00001', 'reason': {'type': 'no_shard_available_action_exception', 'reason': None}}], False), + # Option on: the run continues with the hits from the working shards + (True, [{'index': '.ds-index-test-00001', 'reason': {'type': 'no_shard_available_action_exception', 'reason': None}}], True), + (True, [{'index': '.ds-index-test-00001', 'reason': 'No shard available'}], True), + # Option on: a query that fails to parse stops the run (same behavior as default) + (True, [{'index': '.ds-index-test-00001', 'reason': {'type': 'query_shard_exception', 'reason': 'Failed to parse query'}}], False), +]) +def test_query_with_shard_failures(ea, allow_degraded, failures, expected_success): + ea.allow_queries_on_degraded_indices = allow_degraded + response = generate_hits([START_TIMESTAMP, END_TIMESTAMP]) + response['_shards'] = {'failures': failures} + ea.thread_data.current_es.search.return_value = response + assert ea.run_query(ea.rules[0], START, END) == expected_success + + def test_no_terms_hits(ea): ea.rules[0]['use_terms_query'] = True ea.rules[0]['query_key'] = 'QWERTY' diff --git a/tests/config_test.py b/tests/config_test.py index 5e95a39c..7ee12adb 100644 --- a/tests/config_test.py +++ b/tests/config_test.py @@ -35,6 +35,7 @@ def test_config_loads(): assert conf['writeback_index'] == 'elastalert_status' assert conf['alert_time_limit'] == datetime.timedelta(days=2) + assert conf['allow_queries_on_degraded_indices'] is False def test_config_defaults(): diff --git a/tests/conftest.py b/tests/conftest.py index 0d8c2dcc..3de8e2aa 100644 --- a/tests/conftest.py +++ b/tests/conftest.py @@ -134,6 +134,7 @@ def ea(): 'max_query_size': 10000, 'old_query_limit': datetime.timedelta(weeks=1), 'disable_rules_on_error': False, + 'allow_queries_on_degraded_indices': False, 'scroll_keepalive': '30s', 'custom_pretty_ts_format': '%Y-%m-%d %H:%M'} elastalert.util.elasticsearch_client = mock_es_client From 5682005f71f5142aa7f4825ec72f5796bdf779c5 Mon Sep 17 00:00:00 2001 From: Frank Haldi <177933469+fhaldi@users.noreply.github.com> Date: Fri, 2 Oct 2026 15:28:23 +0200 Subject: [PATCH 5/6] Reduce warning # characters + display shards number --- elastalert/elastalert.py | 8 +++++--- 1 file changed, 5 insertions(+), 3 deletions(-) diff --git a/elastalert/elastalert.py b/elastalert/elastalert.py index fcf8340c..326001a7 100755 --- a/elastalert/elastalert.py +++ b/elastalert/elastalert.py @@ -416,10 +416,12 @@ def get_hits(self, rule, starttime, endtime, index, scroll=False): if parse_errs: raise ElasticsearchException(parse_errs) # If query still worked but returns some errors, continue... - failed_indices = sorted({str(f.get('index')) for f in shard_failures if isinstance(f, dict)}) + failed_indices = ', '.join(sorted({str(f.get('index')) for f in shard_failures if isinstance(f, dict)})) + if len(failed_indices) > 1024: + failed_indices = failed_indices[:1024] + '... (%d characters removed)' % (len(failed_indices) - 1024) elastalert_logger.warning( - 'Rule %s: %d shard(s) failed, continuing with partial results: %s', - rule['name'], len(shard_failures), ', '.join(failed_indices)) + 'Rule %s: %d shard(s) failed, continuing with partial results. Affected indices: %s', + rule['name'], res['_shards'].get('failed', len(shard_failures)), failed_indices) else: # Stop rule if any error with at least one shard of a data stream if len(res.get('_shards', {}).get('failures', [])) > 0: From b5de1fb8c602c045500b65a7b5866890cbdc1f04 Mon Sep 17 00:00:00 2001 From: Frank Haldi <177933469+fhaldi@users.noreply.github.com> Date: Mon, 5 Oct 2026 09:40:05 +0200 Subject: [PATCH 6/6] partial shard failure per rule fixed, test and schema updated --- elastalert/elastalert.py | 2 +- elastalert/schema.yaml | 1 + tests/base_test.py | 13 +++++++++---- 3 files changed, 11 insertions(+), 5 deletions(-) diff --git a/elastalert/elastalert.py b/elastalert/elastalert.py index 326001a7..641b3edb 100755 --- a/elastalert/elastalert.py +++ b/elastalert/elastalert.py @@ -406,7 +406,7 @@ def get_hits(self, rule, starttime, endtime, index, scroll=False): self.thread_data.total_hits = int(res['hits']['total']['value']) # Handles shard failures - if self.allow_queries_on_degraded_indices: + if rule.get('allow_queries_on_degraded_indices', self.allow_queries_on_degraded_indices): # Allow queries on degraded data streams shard_failures = res.get('_shards', {}).get('failures', []) if len(shard_failures) > 0: diff --git a/elastalert/schema.yaml b/elastalert/schema.yaml index f0d6ae9a..bdca5f23 100644 --- a/elastalert/schema.yaml +++ b/elastalert/schema.yaml @@ -302,6 +302,7 @@ properties: query_key: *arrayOfString replace_dots_in_field_names: {type: boolean} scan_entire_timeframe: {type: boolean} + allow_queries_on_degraded_indices: {type: boolean} ## System Error Notifications notify_alert: *stringOrArrayOfStringsOrObjects diff --git a/tests/base_test.py b/tests/base_test.py index 79abdbd7..f5d0c531 100644 --- a/tests/base_test.py +++ b/tests/base_test.py @@ -161,6 +161,7 @@ def test_no_hits(ea): assert ea.rules[0]['type'].add_data.call_count == 0 +@pytest.mark.parametrize('set_in_rule', [False, True]) @pytest.mark.parametrize('allow_degraded, failures, expected_success', [ # All shards working: the option changes nothing (False, [], True), @@ -171,10 +172,14 @@ def test_no_hits(ea): (True, [{'index': '.ds-index-test-00001', 'reason': {'type': 'no_shard_available_action_exception', 'reason': None}}], True), (True, [{'index': '.ds-index-test-00001', 'reason': 'No shard available'}], True), # Option on: a query that fails to parse stops the run (same behavior as default) - (True, [{'index': '.ds-index-test-00001', 'reason': {'type': 'query_shard_exception', 'reason': 'Failed to parse query'}}], False), -]) -def test_query_with_shard_failures(ea, allow_degraded, failures, expected_success): - ea.allow_queries_on_degraded_indices = allow_degraded + (True, [{'index': '.ds-index-test-00001', 'reason': {'type': 'query_shard_exception', 'reason': 'Failed to parse query'}}], False)]) +def test_query_with_shard_failures(ea, set_in_rule, allow_degraded, failures, expected_success): + if set_in_rule: + # The rule's value must win over the global one, so give the global the opposite value + ea.allow_queries_on_degraded_indices = not allow_degraded + ea.rules[0]['allow_queries_on_degraded_indices'] = allow_degraded + else: + ea.allow_queries_on_degraded_indices = allow_degraded response = generate_hits([START_TIMESTAMP, END_TIMESTAMP]) response['_shards'] = {'failures': failures} ea.thread_data.current_es.search.return_value = response