Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
32 changes: 13 additions & 19 deletions scripts/fresh_indices/fresh_indices.py
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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" #############")

Expand Down Expand Up @@ -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" #############")

Expand Down Expand Up @@ -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']={}
Expand All @@ -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']={}
Expand All @@ -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
Expand Down
28 changes: 15 additions & 13 deletions scripts/fresh_indices/fresh_indices.sh
Original file line number Diff line number Diff line change
Expand Up @@ -24,35 +24,37 @@ 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
echo "Before filling the new index, timestamps are captured from the most recently modified"
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"
Expand Down Expand Up @@ -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
Expand Down
11 changes: 6 additions & 5 deletions src/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -24,12 +24,13 @@ python3 -m hubmap_translator <globus-groups-token> 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 <globus-groups-token>" <search-api base URL>/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 <globus-groups-token>" <search-api base URL>/reindex/<uuid>
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 <globus-groups-token>" <search-api base URL>/reindex/<uuid>?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:
Expand Down
Loading
Loading