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
2 changes: 1 addition & 1 deletion README.md
Original file line number Diff line number Diff line change
Expand Up @@ -82,7 +82,7 @@ Make sure to commit newly created migration files.
### Indexing in ES:
- `cd oclapi2/`
- `docker exec -it oclapi2-api-1 python manage.py search_index --populate -f --parallel` -- for populating all indexes
- `docker exec -it oclapi2-api-1 python manage.py search_index --rebuild -f --parallel` -- for rebuild (delete and create) all indexes.
- `docker exec -it oclapi2-api-1 python manage.py search_index --rebuild -f --parallel --use-alias` -- for rebuild all indexes: builds new timestamped indexes, then swaps the aliases to them. Keep `--use-alias` once `concepts`/`mappings` are aliases: without it the rebuild deletes the indexes behind them.
You can also populate/re-index specific indexes, [read more](https://django-elasticsearch-dsl.readthedocs.io/en/latest/management.html)


Expand Down
158 changes: 158 additions & 0 deletions core/common/management/commands/es_split_index.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,158 @@
from datetime import datetime

from django.core.management import BaseCommand, CommandError
from django.utils import timezone
from elasticsearch_dsl.connections import connections

# manage.py es_split_index concepts --shards 6 --dry-run preflight checks and the plan, no changes
# manage.py es_split_index concepts --shards 6 [--yes] split, then swap the alias (asks first unless --yes)
# Hard-links segments into <index>-<timestamp> (writes blocked for minutes), then atomically swaps in an alias <index>.

REQUEST_TIMEOUT = 120


class Command(BaseCommand):
client = None
help = 'Split a single-shard ES index into N primaries behind an alias of the same name.'

def add_arguments(self, parser):
parser.add_argument('index', help='Concrete index to split, e.g. concepts')
parser.add_argument('--shards', type=int, required=True, help='Primary shards of the new index')
parser.add_argument('--dry-run', action='store_true', help='Run the preflight checks only')
parser.add_argument('--yes', action='store_true', help="Don't ask before blocking writes")
parser.add_argument('--timeout', type=int, default=3600, help='Seconds to wait for the new index to be green')
parser.add_argument(
'--max-disk-percent', type=int, default=85,
help="Abort if the shard's node could pass this disk use while the new shards merge apart")

def handle(self, *args, **options):
index = options['index']
shards = options['shards']
self.client = client = connections.get_connection().options(request_timeout=REQUEST_TIMEOUT)

store_bytes, node = self.preflight(index, shards, options['max_disk_percent'])
target = f"{index}-{datetime.now().strftime('%Y%m%d%H%M%S%f')}"
self.stdout.write(
f"Plan: block writes on '{index}' ({store_bytes / 1024 ** 3:.1f} GB on {node}), split it into '{target}' "
f"with {shards} primaries and 0 replicas pinned to {node}, check doc counts, then point alias '{index}' at "
f"'{target}' and delete '{index}'.")
if options['dry_run']:
self.stdout.write('Dry run: nothing changed.')
return
if not options['yes'] and input(f"Block writes on '{index}' and go ahead? [y/N]: ").lower() != 'y':
raise CommandError('Aborted: nothing changed.')

blocked_at = timezone.now().isoformat()
try:
client.indices.add_block(index=index, block='write') # unlike the setting, waits for in-flight writes
self.stdout.write(f"Blocked writes on '{index}' at {blocked_at}.")
self.split(index, target, shards, node, options['timeout'])
self.check_counts(index, target)
# No client retry: a resent swap fails once the first went through, and would look like a failed swap.
client.options(max_retries=0).indices.update_aliases(actions=[
{'add': {'index': target, 'alias': index}},
{'remove_index': {'index': index}},
])
except BaseException as ex: # pylint: disable=broad-exception-caught
swapped = self.is_swapped(index, target)
if not swapped:
self.rollback(index, target, ex, delete_target=swapped is False)
self.stderr.write(f"The swap went through although it reported: {ex}")

self.stdout.write(self.style.SUCCESS(f"Alias '{index}' now points at '{target}'; old index deleted."))
self.stdout.write(
f"If the indexing worker wasn't paused, re-index what changed while writes were blocked: POST "
f"/indexes/resources/{index}/ with filter={{\"updated_at__gte\": \"{blocked_at}\"}}.\n"
f"Next: POST {index}/_forcemerge?only_expunge_deletes=true and wait for merges to finish, then "
f"PUT {index}/_settings {{\"index.routing.allocation.require._name\": null}} to let shards rebalance, "
f"then PUT {index}/_settings {{\"index.number_of_replicas\": 1}}.")

def preflight(self, index, shards, max_disk_percent):
"""Returns (primary store bytes, node) once every check passes; raises CommandError otherwise."""
client = self.client
if shards < 2:
raise CommandError('--shards must be at least 2.')
if client.indices.exists_alias(name=index):
raise CommandError(f"'{index}' is already an alias, nothing to split.")
if not client.indices.exists(index=index):
raise CommandError(f"Index '{index}' doesn't exist.")

index_settings = client.indices.get_settings(index=index)[index]['settings']['index']
if int(index_settings['number_of_shards']) != 1:
raise CommandError(f"'{index}' has {index_settings['number_of_shards']} primaries; only 1 is supported.")
if index_settings.get('blocks'):
raise CommandError(f"'{index}' has blocks set ({index_settings['blocks']}); clear them first.")

health = client.cluster.health(index=index)['status']
if health != 'green':
raise CommandError(f"'{index}' is {health}, not green.")

primary = next(
shard for shard in client.cat.shards(index=index, format='json', bytes='b') if shard['prirep'] == 'p')
store_bytes, node = int(primary['store']), primary['node']
allocation = next(
row for row in client.cat.allocation(format='json', bytes='b') if row['node'] == node)
# Hard-linked source files are only freed once every new shard has merged away from them.
peak_percent = 100 * (int(allocation['disk.used']) + store_bytes) / int(allocation['disk.total'])
if peak_percent > max_disk_percent:
raise CommandError(
f"{node} could reach {peak_percent:.0f}% disk while the new shards merge apart (used + "
f"{store_bytes / 1024 ** 3:.1f} GB), above {max_disk_percent}%.")
return store_bytes, node

def split(self, index, target, shards, node, timeout): # pylint: disable=too-many-arguments
client = self.client
client.indices.flush(index=index)
client.options(request_timeout=timeout).indices.split(
index=index, target=target, settings={
'index.number_of_shards': shards,
'index.number_of_replicas': 0,
'index.blocks.write': None,
# Each new shard holds the whole source via hard links until merged; moving it early copies all that.
'index.routing.allocation.require._name': node,
})
self.stdout.write(f"Split '{index}' into '{target}', waiting for green...")
health = client.options(request_timeout=timeout + 30).cluster.health(
index=target, wait_for_status='green', timeout=f'{timeout}s')
if health['timed_out'] or health['status'] != 'green':
raise CommandError(f"'{target}' isn't green after {timeout}s ({health['status']}).")

def check_counts(self, index, target):
client = self.client
client.indices.refresh(index=[index, target])
source_count = client.count(index=index)['count']
target_count = client.count(index=target)['count']
if source_count != target_count:
raise CommandError(f"Doc counts differ: '{index}' {source_count}, '{target}' {target_count}.")
self.stdout.write(f'Doc counts match: {source_count}.')

def is_swapped(self, index, target):
"""Whether the alias points at the target already; None if ES couldn't be asked."""
try:
return bool(self.client.indices.exists_alias(name=index, index=target))
except Exception: # pylint: disable=broad-exception-caught
return None

def rollback(self, index, target, error, delete_target):
"""Deletes the target (only if it's known not to be serving) and clears the write block, then raises."""
cleanup_errors = []
if delete_target:
try:
self.client.indices.delete(index=target, ignore_unavailable=True)
except Exception as ex: # pylint: disable=broad-exception-caught
cleanup_errors.append(f"deleting '{target}': {ex}")
else:
cleanup_errors.append(f"couldn't tell whether alias '{index}' points at '{target}', so kept it")
try:
self.client.indices.put_settings(index=index, settings={'index.blocks.write': None})
except Exception as ex: # pylint: disable=broad-exception-caught
cleanup_errors.append(f"clearing the write block on '{index}': {ex}")

if cleanup_errors:
raise CommandError(
f"Split failed: {error}. Rollback incomplete ({'; '.join(cleanup_errors)}): check _cat/aliases, "
f"and the write block on '{index}'.") from error
self.stderr.write(f"Rolled back: deleted '{target}' (if created) and cleared the write block on '{index}'.")
if isinstance(error, CommandError):
raise error
raise CommandError(f'Split failed and was rolled back: {error}') from error
29 changes: 18 additions & 11 deletions core/common/models.py
Original file line number Diff line number Diff line change
Expand Up @@ -178,6 +178,15 @@ def attempt(self, start, batch, index_func):
logger.error(message)
ERRBIT_LOGGER.log(BatchIndexingError(message))

def refresh(self, doc):
"""Refreshes the run's index once, if it sent any docs. A failed refresh is only logged."""
if not self.docs:
return
try:
doc._index.refresh() # pylint: disable=protected-access
except Exception as ex: # pylint: disable=broad-except
logger.warning('%s index refresh after batch indexing failed: %s', self.document_name, ex)

def finish(self):
"""Returns the run's summary, or raises BatchIndexingError with it if any batch failed."""
if self.failed_batches:
Expand Down Expand Up @@ -411,8 +420,8 @@ def batch_index_full( # pylint: disable=too-many-arguments
"""
Full (re)index, INDEX_BATCH_SIZE docs a batch (single_batch: all in one). Every batch is attempted; if any
failed, raises BatchIndexingError with the counts once they all have (see BatchIndexRun). Returns the run's
summary.
refresh=False doesn't make ES refresh after each batch (the document's auto_refresh does by default).
summary. refresh (default: the document's auto_refresh) refreshes the index once after the last batch, not
per bulk.
"""
if get(settings, 'TEST_MODE', False):
return None
Expand All @@ -424,15 +433,12 @@ def batch_index_full( # pylint: disable=too-many-arguments
if select_related:
queryset = queryset.select_related(*select_related)

kwargs = {}
refresh = doc.django.auto_refresh if refresh is None else refresh
if refresh: # as django_elasticsearch_dsl's Document.update
kwargs['refresh'] = refresh
run = BatchIndexRun(document)

def index_batch(objects):
run.retry_rejected(lambda: BatchIndexRun.bulk(
doc, doc._get_actions(objects, 'index'), parallel, **kwargs)) # pylint: disable=protected-access
doc, doc._get_actions(objects, 'index'), parallel)) # pylint: disable=protected-access

if single_batch:
run.attempt(0, list(queryset.all()), index_batch)
Expand All @@ -446,6 +452,8 @@ def index_batch(objects):
run.attempt(start, batch, index_batch)
start += batch_size

if refresh:
run.refresh(doc)
return run.finish()

@staticmethod
Expand All @@ -463,15 +471,12 @@ def batch_index_partial_by_ids( # pylint: disable=too-many-arguments
return None

doc = document()
kwargs = {}
refresh = doc.django.auto_refresh if refresh is None else refresh
if refresh:
kwargs['refresh'] = refresh
run = BatchIndexRun(document)

def index_batch(ids):
try:
run.retry_rejected(lambda: BatchIndexRun.bulk(doc, get_actions(ids), parallel, **kwargs))
run.retry_rejected(lambda: BatchIndexRun.bulk(doc, get_actions(ids), parallel))
except BulkIndexError as err:
if on_bulk_error is None:
BaseModel.full_index_missing_docs_or_raise(err, queryset, document)
Expand All @@ -491,6 +496,8 @@ def index_batch(ids):
run.attempt(start, batch, index_batch)
start += batch_size

if refresh:
run.refresh(doc)
return run.finish()

@staticmethod
Expand All @@ -509,7 +516,7 @@ def full_index_missing_docs_or_raise(err, queryset, document, prefetch=None, sel
# Docs not yet in ES -- full index so they appear with all fields
BaseModel.batch_index_full(
single_batch=False, queryset=queryset.filter(id__in={e['update']['_id'] for e in missing}),
document=document, prefetch=prefetch or [], select_related=select_related or []
document=document, prefetch=prefetch or [], select_related=select_related or [], refresh=False
)
missing = []
except BatchIndexingError as ex:
Expand Down
18 changes: 14 additions & 4 deletions core/common/tasks.py
Original file line number Diff line number Diff line change
Expand Up @@ -241,10 +241,20 @@ def __run_search_index_command(command, app_names=None):
if not command:
return

if app_names:
call_command('search_index', f'{command}', '-f', '--models', *app_names, '--parallel')
else:
call_command('search_index', command, '-f', '--parallel')
# Builds a new timestamped index and swaps the alias to it; without this a rebuild breaks an aliased index.
extra_args = ['--use-alias'] if command == '--rebuild' else []
Comment thread
snyaggarwal marked this conversation as resolved.
# --use-alias renames the shared registry Index objects and never renames them back; this process outlives the task.
names = {index: index._name for index in registry.get_indices()} # pylint: disable=protected-access
try:
# refresh=False: no refresh per bulk chunk; the 1s refresh interval makes the docs searchable.
if app_names:
call_command(
'search_index', f'{command}', '-f', '--models', *app_names, '--parallel', *extra_args, refresh=False)
else:
call_command('search_index', command, '-f', '--parallel', *extra_args, refresh=False)
finally:
for index, name in names.items():
index._name = name # pylint: disable=protected-access


@app.task(base=QueueOnceCustomTask, retry_kwargs={'max_retries': 0})
Expand Down
Loading
Loading