From 979fe63bf25499d0aca9867e24ac73c26ecfc20e Mon Sep 17 00:00:00 2001 From: Alec Date: Wed, 24 Jun 2026 17:37:55 -0500 Subject: [PATCH 1/2] using streaming_bulk instead of bulk for more consistent indexing --- condor_history_to_es.py | 9 ++++++--- 1 file changed, 6 insertions(+), 3 deletions(-) diff --git a/condor_history_to_es.py b/condor_history_to_es.py index 7da8585..dd9ae38 100755 --- a/condor_history_to_es.py +++ b/condor_history_to_es.py @@ -52,7 +52,7 @@ def es_generator(entries): yield data from elasticsearch import Elasticsearch -from elasticsearch.helpers import bulk +from elasticsearch.helpers import streaming_bulk prefix = 'http' address = options.address @@ -85,8 +85,11 @@ def es_import(document_generator): json.dump(hit, sys.stdout) success = True else: - success, _ = bulk(es, document_generator, max_retries=20, initial_backoff=10, max_backoff=360) - return success + successes = 0 + for success, _ = streaming_bulk(es, document_generator, max_retries=20, initial_backoff=10, max_backoff=360): + sucesses += success + print("Indexed %d/%d documents" % (successes)) + return successes failed = False if options.access_points and options.collectors: From 03edd43d69c2e94d88bca7282a4d829cd8a7952c Mon Sep 17 00:00:00 2001 From: Alec Date: Thu, 25 Jun 2026 11:17:49 -0500 Subject: [PATCH 2/2] add try except for BulkIndexErrors, fix run_interval field being applied on non-completed jobs --- condor_history_to_es.py | 17 ++++++++++++----- 1 file changed, 12 insertions(+), 5 deletions(-) diff --git a/condor_history_to_es.py b/condor_history_to_es.py index dd9ae38..5f6a984 100755 --- a/condor_history_to_es.py +++ b/condor_history_to_es.py @@ -46,12 +46,14 @@ def es_generator(entries): if options.dailyindex: data['_index'] += '-'+(data['date'].split('T')[0].replace('-','.')) data['_id'] = data['GlobalJobId'].replace('#','-').replace('.','-') - data['run_interval'] = {'gte': data['JobCurrentStartDate'], 'lte': data['EnteredCurrentStatus']} + if data['JobStatus'] == 4: + data['run_interval'] = {'gte': data['JobCurrentStartDate'], 'lte': data['EnteredCurrentStatus']} if not data['_id']: continue yield data from elasticsearch import Elasticsearch +from elasticsearch.helpers import BulkIndexError from elasticsearch.helpers import streaming_bulk prefix = 'http' @@ -75,7 +77,7 @@ def es_generator(entries): es = Elasticsearch(hosts=[url], request_timeout=5000, bearer_auth=token, - sniff_on_connection_fail=True) + sniff_on_node_failure=True) def es_import(document_generator): if options.dry_run: @@ -86,9 +88,14 @@ def es_import(document_generator): success = True else: successes = 0 - for success, _ = streaming_bulk(es, document_generator, max_retries=20, initial_backoff=10, max_backoff=360): - sucesses += success - print("Indexed %d/%d documents" % (successes)) + try: + for success, _ in streaming_bulk(es, document_generator, max_retries=20, initial_backoff=10, max_backoff=360): + successes += success + except BulkIndexError as e: + for error in e.errors: + print(error) + + print(f"Indexed {successes} documents") return successes failed = False