Skip to content
Merged
2 changes: 1 addition & 1 deletion VERSION
Original file line number Diff line number Diff line change
@@ -1 +1 @@
2.3.2
2.3.3
144 changes: 25 additions & 119 deletions metrics/counter/access/accumulation.py
Original file line number Diff line number Diff line change
Expand Up @@ -28,19 +28,10 @@ def accumulate(results, counter_access, line):
access_datetime.date().toordinal(),
access_datetime.hour,
)
compact_accumulate = getattr(results, "accumulate_access", None)
user_session_id = None
if compact_accumulate is None:
user_session_id = _generate_user_session_id(
client_name, client_version, ip_address, access_datetime
)
raw_record = _build_record(
counter_access=counter_access,
line=line,
access_datetime=access_datetime,
second_of_hour=second_of_hour,
user_session_id=user_session_id,
include_id=compact_accumulate is None,
)
access_url_key = access_url or "|".join(
[
Expand All @@ -50,33 +41,18 @@ def accumulate(results, counter_access, line):
]
)

if compact_accumulate is not None:
compact_accumulate(
data=raw_record["data"],
session_key=session_key,
url=access_url_key,
second=second_of_hour,
)
return

item_access_id = raw_record["id"]
if item_access_id not in results:
results[item_access_id] = raw_record["data"]

timestamps_by_url = results[item_access_id].setdefault(
"click_timestamps_by_url", {}
results.accumulate_access(
data=raw_record,
session_key=session_key,
url=access_url_key,
second=second_of_hour,
)
url_timestamps = timestamps_by_url.setdefault(access_url_key, {})
_increment_timestamp_count(url_timestamps, second_of_hour)


def _build_record(
counter_access,
line,
access_datetime,
second_of_hour,
user_session_id,
include_id=True,
):
collection = counter_access.get("collection")
source_key = _source_key(counter_access, collection)
Expand All @@ -89,51 +65,27 @@ def _build_record(
access_country_code = line.get("country_code")
access_date = access_datetime.strftime("%Y-%m-%d")

record = {
"data": {
"collection": collection,
"source_key": source_key,
"document_type": counter_access.get("document_type"),
"pid_v2": pid_v2,
"pid_v3": pid_v3,
"pid_generic": pid_generic,
"document": _document_metadata(counter_access),
"title_pid_generic": counter_access.get("title_pid_generic") or pid_generic,
"user_session_id": user_session_id,
"click_timestamps_by_url": {},
"media_format": media_format,
"content_language": content_language,
"content_type": content_type,
"access_country_code": access_country_code,
"access_date": access_date,
"access_year": access_date[:4],
"access_month": access_date[:7].replace("-", ""),
"publication_year": counter_access.get("publication_year"),
"counter_access_type": counter_access.get("counter_access_type") or "Open",
"access_method": counter_access.get("access_method") or "Regular",
"source": _source_metadata(counter_access),
},
return {
"collection": collection,
"source_key": source_key,
"document_type": counter_access.get("document_type"),
"pid_v2": pid_v2,
"pid_v3": pid_v3,
"pid_generic": pid_generic,
"document": _document_metadata(counter_access),
"title_pid_generic": counter_access.get("title_pid_generic") or pid_generic,
"media_format": media_format,
"content_language": content_language,
"content_type": content_type,
"access_country_code": access_country_code,
"access_date": access_date,
"access_year": access_date[:4],
"access_month": access_date[:7].replace("-", ""),
"publication_year": counter_access.get("publication_year"),
"counter_access_type": counter_access.get("counter_access_type") or "Open",
"access_method": counter_access.get("access_method") or "Regular",
"source": _source_metadata(counter_access),
}
if include_id:
record["id"] = _generate_item_access_id(
user_session_id=user_session_id,
col_acron3=collection,
source_key=source_key,
pid_v2=pid_v2,
pid_v3=pid_v3,
pid_generic=pid_generic,
content_language=content_language,
access_country_code=access_country_code,
media_format=media_format,
content_type=content_type,
)
return record


def _increment_timestamp_count(timestamps, key):
if key not in timestamps:
timestamps[key] = 0
timestamps[key] += 1


def _normalized_access_path(url):
Expand All @@ -150,23 +102,6 @@ def _normalized_access_path(url):
return path or None


def _generate_user_session_id(
client_name, client_version, ip_address, datetime, sep="|"
):
dt_year_month_day = datetime.strftime("%Y-%m-%d")
dt_hour = datetime.strftime("%H")

return sep.join(
[
str(client_name),
str(client_version),
str(ip_address),
str(dt_year_month_day),
str(dt_hour),
]
)


def _document_metadata(counter_access):
document_title = counter_access.get("document_title")
return {"title": document_title} if document_title else {}
Expand Down Expand Up @@ -196,32 +131,3 @@ def _source_key(counter_access, fallback):
or counter_access.get("source_type")
or fallback
)


def _generate_item_access_id(
col_acron3,
source_key,
pid_v2,
pid_v3,
pid_generic,
user_session_id,
access_country_code,
content_language,
media_format,
content_type,
sep="|",
):
return sep.join(
[
col_acron3,
str(source_key or ""),
pid_v2 or "",
pid_v3 or "",
pid_generic or "",
str(user_session_id or ""),
str(access_country_code or ""),
str(content_language or ""),
str(media_format or ""),
str(content_type or ""),
]
)
65 changes: 27 additions & 38 deletions metrics/counter/access/daily_accumulator.py
Original file line number Diff line number Diff line change
Expand Up @@ -117,11 +117,11 @@ def _timestamps_as_dict(self, accumulator):
return timestamps


class DailyAccessAccumulator(dict):
class DailyAccessAccumulator:
"""Store compact records and materialize them only for metric conversion."""

def __init__(self):
super().__init__()
self._records = {}
self._documents = [None, {}]
self._document_ids = {}
self._sources = [None]
Expand All @@ -130,28 +130,8 @@ def __init__(self):
self._strings = [None]
self._string_ids = {}

def __setitem__(self, key, value):
if isinstance(value, _CompactAccessRecord):
super().__setitem__(key, value)
return

source_key = value.get("source_key")
source = value.get("source")
if source_key and source:
source_id = self._intern_source(source_key, source)
value["source"] = self._sources[source_id]

document_key = self._document_key(value)
document = value.get("document")
if document is not None and any(document_key[1:]):
document_id = self._intern_document(value)
value["document"] = self._documents[document_id]

user_session_id = value.get("user_session_id")
if user_session_id:
value["user_session_id"] = self._legacy_intern_session(user_session_id)

super().__setitem__(key, value)
def __len__(self):
return len(self._records)

def accumulate_access(self, data, session_key, url, second):
session = self._intern_session(session_key)
Expand All @@ -167,19 +147,35 @@ def accumulate_access(self, data, session_key, url, second):
self._intern(data.get("media_format")),
self._intern(data.get("content_type")),
)
record = dict.get(self, key)
record = self._records.get(key)
if record is None:
record = _CompactAccessRecord(self, data, session, url, second)
dict.__setitem__(self, key, record)
self._records[key] = record
return
record.add_timestamp(self._intern(url), second)

def iter_materialized_values(self):
for value in dict.values(self):
if isinstance(value, _CompactAccessRecord):
def iter_materialized_values(self, consume=False):
if not consume:
for value in self._records.values():
yield value.as_dict(self)
else:
yield value
return

keys = tuple(self._records)
try:
for key in keys:
yield self._records.pop(key).as_dict(self)
finally:
self.clear()

def clear(self):
self._records.clear()
self._documents.clear()
self._document_ids.clear()
self._sources.clear()
self._source_ids.clear()
self._sessions.clear()
self._strings.clear()
self._string_ids.clear()

def _intern(self, value):
if value is None:
Expand Down Expand Up @@ -208,13 +204,6 @@ def _intern_session(self, session_key):
self._sessions[compact_key] = session_id
return session_id

def _legacy_intern_session(self, session):
interned = self._sessions.get(session)
if interned is None:
self._sessions[session] = session
return session
return interned

def _intern_source(self, source_key, source):
if source is None:
return _NONE_METADATA_ID
Expand Down
15 changes: 2 additions & 13 deletions metrics/counter/indexing/converter.py
Original file line number Diff line number Diff line change
Expand Up @@ -14,22 +14,11 @@
_DEFAULT = DocumentPipeline()


def convert(data):
if not isinstance(data, dict):
return {"month": {}, "year": {}}

month_data = _convert_granularity(data, "month")
year_data = _convert_granularity(data, "year")

return {"month": month_data, "year": year_data}


def _convert_granularity(data, granularity):
def convert_granularity(values, granularity):
converted_data = {}
unique_state = _initialize_unique_state()
values = getattr(data, "iter_materialized_values", data.values)

for value in values():
for value in values:
pipeline = _get_pipeline(value)
pipeline.accumulate(
data=converted_data,
Expand Down
30 changes: 21 additions & 9 deletions metrics/opensearch/client.py
Original file line number Diff line number Diff line change
Expand Up @@ -10,10 +10,13 @@
merge_metric_document,
)

_BULK_CHUNK_SIZE = 500


class OpenSearchUsageClient:
def __init__(self, url=None, basic_auth=None, api_key=None, verify_certs=None):
self.client = self.get_opensearch_client(url, basic_auth, api_key, verify_certs)
logging.info("OpenSearch HTTP request compression is enabled.")

def get_opensearch_client(
self,
Expand All @@ -33,10 +36,20 @@ def get_opensearch_client(
url,
http_auth=tuple(basic_auth),
verify_certs=verify_certs,
http_compress=True,
)
if api_key:
return OpenSearch(url, api_key=api_key, verify_certs=verify_certs)
return OpenSearch(url, verify_certs=verify_certs)
return OpenSearch(
url,
api_key=api_key,
verify_certs=verify_certs,
http_compress=True,
)
return OpenSearch(
url,
verify_certs=verify_certs,
http_compress=True,
)

def ping(self):
try:
Expand Down Expand Up @@ -104,20 +117,17 @@ def index_documents(self, index_name, documents, ping_client=False):
),
)

def increment_documents_for_daily_job(
def increment_document_items_for_daily_job(
self,
index_name,
documents,
document_items,
job_id,
ping_client=False,
):
if ping_client and not self.ping():
return

if not documents:
return

helpers.bulk(
succeeded, _failed = helpers.bulk(
self.client,
(
build_idempotent_job_increment_action(
Expand All @@ -126,9 +136,11 @@ def increment_documents_for_daily_job(
document=document,
job_id=job_id,
)
for doc_id, document in documents.items()
for doc_id, document in document_items
),
chunk_size=_BULK_CHUNK_SIZE,
)
return succeeded

def delete_documents(self, index_name, doc_ids, ping_client=False):
if ping_client and not self.ping():
Expand Down
Loading
Loading