diff --git a/scripts/fresh_indices/fresh_indices.py b/scripts/fresh_indices/fresh_indices.py index 3f29dc99..19cef3fd 100644 --- a/scripts/fresh_indices/fresh_indices.py +++ b/scripts/fresh_indices/fresh_indices.py @@ -363,7 +363,7 @@ def swap_index_names_per_strategy(es_mgr:ESManager, fill_strategy:FillStrategyTy destination_index=index_info_dict[source_index]['destination'] flush_index=destination_index.replace('fill','flush') - # Block writing on the indices, even though services which write to them should probably be down. + # Block writing on the indices temporarily during the index swap operation. logger.debug(f"Set {IndexBlockType.WRITE.value} block on source_index={source_index}.") es_mgr.set_index_block(index_name=source_index , block_type_enum=IndexBlockType.WRITE) @@ -430,10 +430,9 @@ def swap_index_names_per_strategy(es_mgr:ESManager, fill_strategy:FillStrategyTy f" #############") -# Read the op_data file from the last 'create' command. Read each document from -# the flush index which was created or updated after the 'create' command started. -# Re-index those entities into the new index, even though re-indexing is a more -# expensive operations, so that the new index has everything the flush index had. +# Read the op_data file from the last 'create' command. Identify documents in the flush index +# which were modified after the 'create' command started, and enqueue them for re-indexing +# into the now-live indices via background workers. def catch_up_live_index(es_mgr:ESManager)->None: global op_data global op_data_supplement @@ -534,10 +533,9 @@ def catch_up_live_index(es_mgr:ESManager)->None: end_time = time.time() # KBKBKB @TODO check in with Joe if it is worth it to try determining if threads err'ed and pointing that out here... - logger.info(f"############# Re-indexing entities of recently touch documents via script complete at {time.strftime('%H:%M:%S',time.localtime(end_time))} #############") - + logger.info(f"############# Re-indexing entities of recently touched documents enqueuing complete at {time.strftime('%H:%M:%S',time.localtime(end_time))} #############") elapsed_seconds = end_time-start_time - logger.info(f"############# Re-indexing via script took" + logger.info(f"############# Re-indexing entity enqueuing took" f" {time.strftime('%H:%M:%S', time.gmtime(elapsed_seconds))}." f" #############") @@ -571,10 +569,9 @@ def create_new_indices(): end_time = time.time() # KBKBKB @TODO check in with Joe if it is worth it to try determining if threads err'ed and pointing that out here... - logger.info(f"############# Full index via script complete at {time.strftime('%H:%M:%S',time.localtime(end_time))} #############") - + logger.info(f"############# Full index job enqueuing complete at {time.strftime('%H:%M:%S',time.localtime(end_time))} #############") elapsed_seconds = end_time-start_time - logger.info(f"############# Full index via script took" + logger.info(f"############# Full index job enqueuing took" f" {time.strftime('%H:%M:%S', time.gmtime(elapsed_seconds))}." f" #############") @@ -657,10 +654,8 @@ def create_new_indices(): create_new_indices() print('#############') print('Completed create command.') - print(f"Next either take down the service and execute the 'catch-up'" - f" command, or execute 'catch-up' with the service up, knowing" - f" is may need to be executed again if documents are being" - f" written to the indices right now.") + print(f"All entities have been enqueued for indexing. Wait for workers to finish" + f" filling the indices before executing the 'catch-up' command.") print('#############') elif command == 'catch-up': op_data_supplement['catchup']={} @@ -669,9 +664,8 @@ def create_new_indices(): catch_up_live_index(es_mgr=esmanager) print('#############') print('Completed catch-up command.') - print(f"Next either take down the service and catch-up again before" - f" executing 'go-live', or proceed to execute 'go-live' if" - f" no new documents are being written to the indices right now.") + print(f"Catch-up entities have been enqueued for indexing via background workers." + f" Workers will re-index recently modified entities into the live indices.") print('#############') elif command == 'go-live': op_data_supplement['golive']={} @@ -681,7 +675,7 @@ def create_new_indices(): print('#############') print('Completed go-live command.') print(f"You may want to visually verify all ElasticSearch indices have 'green health' in AWS." - f" If services were brought down to execute this command, they can be brought back up.") + f" Run 'catch-up' next to re-index any entities modified since the 'create' command started.") print('#############') # If executing the preceding commands generated any extra operational data to be externalized for diff --git a/scripts/fresh_indices/fresh_indices.sh b/scripts/fresh_indices/fresh_indices.sh index 9f652b40..54ba5a1a 100755 --- a/scripts/fresh_indices/fresh_indices.sh +++ b/scripts/fresh_indices/fresh_indices.sh @@ -24,7 +24,7 @@ Help() echo "Script to initialize a fresh ElasticSearch index using Neo4j data." echo echo "A new ElasticSearch index is created using the strategy specified in the fresh_index.ini file." - echo "The new index is created and mostly filled without taking the search-api offline. During" + echo "The new index is created and filled without taking the search-api offline. During" echo "this time, the new index has a temporary name, and the existing index continues supporting" echo "Production." echo @@ -32,27 +32,29 @@ Help() echo "documents in the ElasticSearch index with every document. This index is configured in" echo "search-config.yaml as 'entities'->default_index->'private' i.e. typically hm_*_consortium_entities." echo - echo "After the new index is filled, the search-api service must be taken offline. Then, any documents" - echo "modified after the captured timestamps are 're-indexed' in the new index. After those documents" - echo "are refreshed to reflect activity that happened after the timestamps were captured, the existing" - echo "index is renamed and can be manually deleted. The new index is renamed to support Production." + echo "The 'create' command enqueues all indexing jobs to background workers via a Redis job queue." + echo "The script returns once enqueuing is complete; workers fill the new index in the background." + echo "Monitor the reindex-status endpoint to determine when workers have finished." + echo + echo "Once workers have finished, run 'go-live' to swap the new index into production." + echo "Then run 'catch-up' to re-enqueue any entities that were modified during the fill period." + echo "Workers will process the catch-up jobs against the now-live indices. The search-api" + echo "does not need to be taken offline at any point in this workflow." echo echo "When the index renaming activity is compete, this script will wait for the ElasticSearch index health to" echo "become green. After that, the search-api can be returned to service, and will use the fresh index." echo echo "Syntax: $0 [-option] [command]" echo "[command]" - echo "create - Create a new ElasticSearch index with a temporary name, filled with documents indexed from Neo4j data." - echo "catch-up - While search-api is down" - echo " * re-index documents which were modified after the create started," - echo " * rename the old index so it can be deleted," - echo " * rename the new index for use by Production," - echo " * and wait for the new index to have green health." - echo "go-live - Swap indices around so the results of 'create' and 'catch-up' commands becomes the indices used by the Search API." + echo "create - Create new ElasticSearch indices with temporary names, and enqueue all indexing jobs" + echo " to background workers. Returns when enqueuing is complete; workers fill indices in background." + echo "go-live - Swap the filled temporary indices into production once workers have finished filling them." echo " * Names of current and new indices are taken from the newest exec_info/op_data*.json file." echo " * The current indices will be renamed with a 'flush' prefix." echo " * The new indices will taken on the name expected by Search API." echo " * The script will wait for 'green health' on each renamed index." + echo "catch-up - After go-live, identify any entities modified since the 'create' command started" + echo " and enqueue them for re-indexing against the now-live indices via background workers." echo echo "[-option]" echo "-h Display this help" @@ -248,7 +250,7 @@ if [[ "$cmd" == "create" ]]; then echo -e "\t$index " done elif [[ "$cmd" == "catch-up" ]]; then - echo "Using op_data from the most current $arg_output_dir/op_data*.json file to re-index any entities touch since the 'create' command." + echo "Using op_data from the most current $arg_output_dir/op_data*.json file to enqueue re-indexing of any entities modified since the 'create' command." elif [[ "$cmd" == "go-live" ]]; then echo "Using op_data from the most current $arg_output_dir/op_data*.json file, swapping index names so Search API can use new indices, " else diff --git a/src/README.md b/src/README.md index baacddd9..bb9ab65c 100644 --- a/src/README.md +++ b/src/README.md @@ -24,12 +24,13 @@ python3 -m hubmap_translator 1>indexer.log 2>&1 The live reindex will NOT recreate the indices, instead it will just delete and documents that are no longer in Neo4j and reindex each entity document found in Neo4j. -```` -curl -i -X PUT -H "Authorization:Bearer " /reindex-all -```` - -The token will need to be in the admin group. +Individual entity reindex requests are handled via a Redis-backed priority queue. When a reindex request is received, the target entity and all of its related entities (ancestors, descendants, revisions, collections, uploads) are enqueued as separate jobs rather than executed immediately. Jobs are processed by worker processes defined in `jobq_workers.py`, which run in the same container as the main service but are started independently of the uWSGI process. Redis must be running and the worker processes must be active for reindex jobs to be executed. +To reindex a single entity: +curl -i -X PUT -H "Authorization:Bearer " /reindex/ +An optional `priority` query parameter controls job priority. Valid values are `1`, `2`, and `3`, where `1` is the highest priority and the default. When a job is enqueued at priority `1`, its related entities are enqueued at priority `2`. Jobs enqueued at priority `2` or `3` have their related entities enqueued at the same priority level. +curl -i -X PUT -H "Authorization:Bearer " /reindex/?priority=2 +To reindex all entities, use the scripts found within `scripts/fresh_indices`. See the `readme.me` in that directory for more details. This replaces the endpoint `/reindex-all`. ## To debug Capture one or more documents which fail during indexing. Then, from `src/` run: diff --git a/src/hubmap_translator.py b/src/hubmap_translator.py index be77d929..9faf62c3 100644 --- a/src/hubmap_translator.py +++ b/src/hubmap_translator.py @@ -1,5 +1,6 @@ import concurrent.futures import copy +import contextvars import importlib import requests import json @@ -31,7 +32,14 @@ from translator.translator_interface import TranslatorInterface logger = logging.getLogger(__name__) +job_context = contextvars.ContextVar('job_context', default='LIVE') +class JobContextFilter(logging.Filter): + def filter(self, record): + record.msg = f"[{job_context.get()}] {record.msg}" + return True + +logger.addFilter(JobContextFilter()) config = {} app = Flask(__name__, instance_path=os.path.join(os.path.abspath(os.path.dirname(__file__)), 'instance'), instance_relative_config=True) @@ -225,74 +233,6 @@ def log_configuration(self, log_level:int=logger.getEffectiveLevel()): logger.log(level=log_level , msg=f"\tTRANSFORMERS={self.TRANSFORMERS}") - # Used by full reindex via script and live reindex-all call - def translate_all(self): - with app.app_context(): - try: - logger.info("Start executing translate_all()") - - start = time.time() - - donor_uuids_list = get_uuids_by_entity_type("donor", self.request_headers, self.DEFAULT_ENTITY_API_URL) - upload_uuids_list = get_uuids_by_entity_type("upload", self.request_headers, self.DEFAULT_ENTITY_API_URL) - collection_uuids_list = get_uuids_by_entity_type("collection", self.request_headers, self.DEFAULT_ENTITY_API_URL) - - # Only need this comparision for the live /rindex-all PUT call - if not self.skip_comparision: - # Make calls to entity-api to get a list of uuids for rest of entity types - sample_uuids_list = get_uuids_by_entity_type("sample", self.request_headers, self.DEFAULT_ENTITY_API_URL) - dataset_uuids_list = get_uuids_by_entity_type("dataset", self.request_headers, self.DEFAULT_ENTITY_API_URL) - - # Merge into a big list that with no duplicates - all_entities_uuids = set(donor_uuids_list + sample_uuids_list + dataset_uuids_list + upload_uuids_list + collection_uuids_list) - - es_uuids = [] - index_names = get_all_reindex_enabled_indice_names(self.INDICES) - - for index in index_names.keys(): - all_indices = index_names[index] - # get URL for that index - es_url = self.INDICES['indices'][index]['elasticsearch']['url'].strip('/') - - for actual_index in all_indices: - es_uuids.extend(get_uuids_from_es(actual_index, es_url)) - - es_uuids = set(es_uuids) - - # Remove entities found in Elasticsearch but no longer in neo4j - for uuid in es_uuids: - if uuid not in all_entities_uuids: - logger.debug(f"Entity of uuid: {uuid} found in Elasticsearch but no longer in neo4j. Delete it from Elasticsearch.") - self.delete(uuid) - - with concurrent.futures.ThreadPoolExecutor() as executor: - # The default number of threads in the ThreadPoolExecutor is calculated as: - # From 3.8 onwards default value is min(32, os.cpu_count() + 4) - # Where the number of CPUs is determined by Python and will take hyperthreading into account - logger.info(f"The number of worker threads being used by default: {executor._max_workers}") - - # Submit tasks to the thread pool - collection_futures_list = [executor.submit(self.translate_collection, uuid, reindex=True) for uuid in collection_uuids_list] - upload_futures_list = [executor.submit(self.translate_upload, uuid, reindex=True) for uuid in upload_uuids_list] - - # Append the above lists into one - futures_list = collection_futures_list + upload_futures_list - - # The target function runs the task logs more details when f.result() gets executed - for f in concurrent.futures.as_completed(futures_list): - result = f.result() - - # Index the donor tree in a regular for loop, not the concurrent mode - # However, the descendants of a given donor will be indexed concurrently - for uuid in donor_uuids_list: - self.translate_donor_tree(uuid) - - end = time.time() - - logger.info(f"Finished executing translate_all(). Total time used: {end - start} seconds.") - except Exception as e: - logger.error(e) - # Used by full reindex scripts only. # Assumes the index named indices are already created and empty. # Require Data Admin privileges to execute. @@ -329,7 +269,7 @@ def translate_full(self, reindex_queue=None, index_override=None): if reindex_queue is not None: all_uuids = donor_uuids_list + upload_uuids_list + collection_uuids_list for uuid in all_uuids: - self.enqueue_reindex(uuid, reindex_queue, priority=1, index_override=index_override) + self.enqueue_reindex(uuid, reindex_queue, priority=1, index_override=index_override, job_type='FULL') else: with concurrent.futures.ThreadPoolExecutor() as executor: logger.info(f"The number of worker threads being used by default: {executor._max_workers}") @@ -675,13 +615,14 @@ def _transform_and_write_entity_to_index_group(self, entity: dict, index_group: f" entity['entity_type']={entity['entity_type']}") - def enqueue_reindex(self, entity_id, reindex_queue, priority, index_override=None): + def enqueue_reindex(self, entity_id, reindex_queue, priority, index_override=None, job_type='LIVE'): + job_context.set(job_type) try: logger.info(f"Start executing translate() on entity_id: {entity_id}") entity = self.call_entity_api(entity_id=entity_id, endpoint_base='documents') logger.info(f"Enqueueing reindex for {entity['entity_type']} of uuid: {entity_id}") subsequent_priority = max(priority, 2) - kwargs_for_job = {} + kwargs_for_job = {'job_type': job_type} if index_override: kwargs_for_job['index_override'] = index_override reference_id = reindex_queue.enqueue( @@ -784,7 +725,7 @@ def enqueue_reindex(self, entity_id, reindex_queue, priority, index_override=Non if response.status_code == 200: associated_metadata = response.json() else: - self.logger.error(f"Failed to fetch batch metadata: {response.status_code}") + logger.error(f"Failed to fetch batch metadata: {response.status_code}") associated_metadata = {} except Exception as e: logger.error(f"Unable to retrieve uuid and hubmap_id from entity-api. Proceed with enqueuing but this info will be missing from logging and status. {e}") @@ -794,7 +735,7 @@ def enqueue_reindex(self, entity_id, reindex_queue, priority, index_override=Non jobs.append({ "entity_id": related_entity_id, "args": [related_entity_id, self.token], - "kwargs": {"index_override": index_override} if index_override else {}, + "kwargs": {"index_override": index_override, "job_type": job_type} if index_override else {"job_type": job_type}, "metadata": meta, }) if jobs: @@ -1358,13 +1299,13 @@ def load_public_doc_exclusion_dict(self, entity_api_prov_schema_raw_url): raise YAMLError(ye) else: msg = f"Unable to retrieve public index field exclusion information" - self.logger.error( f"{msg}." + logger.error( f"{msg}." f" Got an HTTP {response.status_code}" f" retrieving {self.indices['entity_api_prov_schema_raw_url']}") raise HTTPException(f"{msg}. See logs.") if not provenance_schema_dict or 'ENTITIES' not in provenance_schema_dict: msg = f"Unable retrieve Entity API's provenance_schema.yaml information" - self.logger.error( f"{msg}." + logger.error( f"{msg}." f" Not expected content using the translator's" f" self.indices['entity_api_prov_schema_raw_url']={self.indices['entity_api_prov_schema_raw_url']}.") raise Exception(f"{msg}. See logs.") @@ -1825,8 +1766,12 @@ def _generate_public_doc(self, entity, index_group:str): # Can't reuse call_entity_api() here due to the response data type # Making a call against entity-api/entities/?property=status url = self.entity_api_url + "/entities/" + next_revision_uuid + "?property=status" - response = requests.get(url, headers=self.request_headers, verify=False) - + try: + response = requests.get(url, headers=self.request_headers, verify=False) + except Exception as e: + msg = f"_generate_public_doc() failed to get Dataset/Publication status of next_revision_uuid via entity-api for uuid: {next_revision_uuid}" + logger.error(msg) + raise Exception(msg) if response.status_code != 200: logger.error(f"_generate_public_doc() failed to get Dataset/Publication status of next_revision_uuid via entity-api for uuid: {next_revision_uuid}") @@ -1942,9 +1887,19 @@ def call_entity_api(self, entity_id, endpoint_base, endpoint_suffix=None, url_pr url = f"{url}/{endpoint_suffix}" if url_property: url = f"{url}?property={url_property}" + try: + response = requests.get(url, headers=self.request_headers, verify=False) + except Exception as e: + msg = f"call_entity_api() failed to get entity of uuid {entity_id} via entity-api" + logger.exception(msg) + # Add this uuid to the failed list + self.failed_entity_api_calls.append(url) + self.failed_entity_ids.append(entity_id) - response = requests.get(url, headers=self.request_headers, verify=False) - + # Bubble up the error message from entity-api instead of sys.exit(msg) + # The caller will need to handle this exception + raise requests.exceptions.RequestException(msg) + if response.status_code != 200: msg = f"call_entity_api() failed to get entity of uuid {entity_id} via entity-api" @@ -1977,8 +1932,13 @@ def get_collection_doc(self, entity_id): # - a valid token but not in HuBMAP-Read group or # - no token at all # Here we do NOT send over the token - url = self.entity_api_url + "/documents/" + entity_id - response = requests.get(url, headers=self.request_headers, verify=False) + try: + url = self.entity_api_url + "/documents/" + entity_id + response = requests.get(url, headers=self.request_headers, verify=False) + except Exception as e: + msg = f"get_collection_doc() failed to get entity of uuid {entity_id} via entity-api" + logger.error(msg) + raise if response.status_code != 200: msg = f"get_collection_doc() failed to get entity of uuid {entity_id} via entity-api" @@ -2097,7 +2057,8 @@ def get_organ_types(self): # This approach is different from the live /reindex-all PUT call # It'll delete all the existing indices and recreate then then index everything -def reindex_entity_queued_wrapper(entity_id, token, index_override=None): +def reindex_entity_queued_wrapper(entity_id, token, index_override=None, job_type='LIVE'): + job_context.set(job_type) indices = index_override if index_override else config['INDICES'] translator = Translator( indices=indices,