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
1 change: 1 addition & 0 deletions config/collections.py
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@
"preprints": "small",
"pry": "small",
"rve": "small",
"rvt": "small",
"spa": "small",
"sss": "small",
"sza": "small",
Expand Down
4 changes: 4 additions & 0 deletions config/settings/base.py
Original file line number Diff line number Diff line change
Expand Up @@ -488,3 +488,7 @@
# Collection size categories
# ------------------------------------------------------------------------------
SUPPORTED_LOGFILE_EXTENSIONS = env.list("SUPPORTED_LOGFILE_EXTENSIONS", default=[".log", ".gz", ".zip"])
PARSING_METADATA_CACHE_COLLECTIONS = env.list(
"PARSING_METADATA_CACHE_COLLECTIONS",
default=list(COLLECTION_ACRON3_SIZE_MAP),
)
29 changes: 21 additions & 8 deletions metrics/services/parsing/metadata.py
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@

from document.models import Document
from log_manager_config.models import CollectionLogDirectory
from metrics.services.parsing import metadata_cache
from source.models import Source

TRANSLATOR_CLASSES = {
Expand All @@ -30,20 +31,32 @@ def build_url_translation_manager(log_file):
f"No URL translator class found for collection {log_file.collection}."
)

started = monotonic()
manager = url_translator.URLTranslationManager(
documents_metadata=Document.metadata(collection=log_file.collection),
sources_metadata=Source.metadata(collection=log_file.collection),
translator=translator_class,
)
if metadata_cache.is_enabled(log_file.collection):
return metadata_cache.get_url_translation_manager(
log_file.collection,
translator_class,
_build_manager,
)

manager, elapsed = _build_manager(log_file.collection, translator_class)
logging.info(
"Prepared parsing metadata for %s in %.3f seconds.",
"Prepared parsing metadata for %s without cache in %.3f seconds.",
log_file.collection.acron3,
monotonic() - started,
elapsed,
)
return manager


def _build_manager(collection, translator_class):
started = monotonic()
manager = url_translator.URLTranslationManager(
documents_metadata=Document.metadata(collection=collection),
sources_metadata=Source.metadata(collection=collection),
translator=translator_class,
)
return manager, monotonic() - started


def _get_log_file_translator_class(log_file):
for directory in CollectionLogDirectory.objects.filter(
config__collection=log_file.collection,
Expand Down
144 changes: 144 additions & 0 deletions metrics/services/parsing/metadata_cache.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,144 @@
import logging
from copy import copy
from threading import Lock
from time import monotonic

from django.conf import settings
from django.db import connection, transaction
from django.db.models import Count, Max

from config.collections import get_collection_size
from document.models import Document
from source.models import Source

_CACHE_ENTRY = None
_CACHE_LOCK = Lock()


def is_enabled(collection):
enabled_collections = {
value.lower()
for value in getattr(settings, "PARSING_METADATA_CACHE_COLLECTIONS", [])
}
return collection.acron3.lower() in enabled_collections


def get_url_translation_manager(collection, translator_class, build_manager):
global _CACHE_ENTRY

acronym = collection.acron3
size = get_collection_size(acronym)
translator_name = translator_class.__name__

with _CACHE_LOCK:
entry = _CACHE_ENTRY
reason = _get_rebuild_reason(entry, collection, translator_name)

if reason is None:
signature = _read_signature(collection)
if signature == entry["signature"]:
logging.info(
"Parsing metadata cache hit for %s (size=%s, signature=%s).",
acronym,
size,
signature,
)
return _fresh_manager(entry)
reason = "signature_changed"

started = monotonic()
new_entry = _build_cache_entry(
collection,
translator_class,
build_manager,
)
elapsed = monotonic() - started
_CACHE_ENTRY = new_entry

logging.info(
"Parsing metadata cache %s for %s "
"(size=%s, reason=%s, signature=%s, build_seconds=%.3f).",
"miss" if entry is None else "rebuild",
acronym,
size,
reason,
new_entry["signature"],
elapsed,
)
return _fresh_manager(new_entry)


def clear():
global _CACHE_ENTRY

with _CACHE_LOCK:
_CACHE_ENTRY = None


def _get_rebuild_reason(entry, collection, translator_name):
if entry is None:
return "empty"
if entry["collection_id"] != collection.pk:
return "collection_changed"
if entry["translator_name"] != translator_name:
return "translator_changed"
return None


def _build_cache_entry(collection, translator_class, build_manager):
if connection.in_atomic_block:
return _load_cache_entry(collection, translator_class, build_manager)

with transaction.atomic():
with connection.cursor() as cursor:
cursor.execute("SET TRANSACTION ISOLATION LEVEL REPEATABLE READ, READ ONLY")
return _load_cache_entry(collection, translator_class, build_manager)


def _load_cache_entry(collection, translator_class, build_manager):
signature = _get_collection_signature(collection)
manager, _ = build_manager(collection, translator_class)
return {
"collection_id": collection.pk,
"translator_class": translator_class,
"translator_name": translator_class.__name__,
"signature": signature,
"manager": manager,
}


def _read_signature(collection):
if connection.in_atomic_block:
return _get_collection_signature(collection)

with transaction.atomic():
with connection.cursor() as cursor:
cursor.execute("SET TRANSACTION ISOLATION LEVEL REPEATABLE READ, READ ONLY")
return _get_collection_signature(collection)


def _get_collection_signature(collection):
document_signature = Document.objects.filter(collection=collection).aggregate(
count=Count("pk"),
max_updated=Max("updated"),
)
source_signature = Source.objects.filter(collection=collection).aggregate(
count=Count("pk"),
max_updated=Max("updated"),
)
return (
document_signature["count"],
document_signature["max_updated"],
source_signature["count"],
source_signature["max_updated"],
)


def _fresh_manager(entry):
manager = copy(entry["manager"])
manager.translator = entry["translator_class"](
manager.sources_metadata,
manager.documents_metadata,
)
manager.is_translator_forced = True
return manager
3 changes: 2 additions & 1 deletion metrics/tests/parsing/test_metadata.py
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
from types import SimpleNamespace
from unittest.mock import patch

from django.test import TestCase
from django.test import TestCase, override_settings

from collection.models import Collection
from log_manager_config.models import CollectionLogDirectory, LogManagerCollectionConfig
Expand Down Expand Up @@ -57,6 +57,7 @@ def setUp(self):
translator_class="books",
)

@override_settings(PARSING_METADATA_CACHE_COLLECTIONS=[])
@patch("metrics.services.parsing.metadata.url_translator.URLTranslationManager")
@patch("metrics.services.parsing.metadata.Source.metadata")
@patch("metrics.services.parsing.metadata.Document.metadata")
Expand Down
Loading
Loading