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 b9a6be93..641b3edb 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', []) @@ -404,14 +405,33 @@ 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 + 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: + 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 = ', '.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. 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: + 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/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/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' } diff --git a/tests/base_test.py b/tests/base_test.py index 6a345b61..f5d0c531 100644 --- a/tests/base_test.py +++ b/tests/base_test.py @@ -161,6 +161,31 @@ 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), + (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, 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 + 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