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
60 changes: 40 additions & 20 deletions metrics/counter/access/accumulation.py
Original file line number Diff line number Diff line change
Expand Up @@ -21,31 +21,48 @@ def accumulate(results, counter_access, line):
access_datetime = local_datetime.replace(minute=0, second=0, microsecond=0)
second_of_hour = local_datetime.minute * 60 + local_datetime.second

user_session_id = _generate_user_session_id(
session_key = (
client_name,
client_version,
ip_address,
access_datetime,
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,
)
item_access_id = raw_record["id"]

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

access_url_key = access_url or "|".join(
[
str(counter_access.get("pid_generic") or ""),
str(counter_access.get("media_format") or ""),
str(counter_access.get("content_type") or ""),
]
)

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", {}
)
Expand All @@ -59,6 +76,7 @@ def _build_record(
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 @@ -71,19 +89,7 @@ def _build_record(
access_country_code = line.get("country_code")
access_date = access_datetime.strftime("%Y-%m-%d")

return {
"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,
),
record = {
"data": {
"collection": collection,
"source_key": source_key,
Expand All @@ -108,6 +114,20 @@ def _build_record(
"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):
Expand Down
259 changes: 241 additions & 18 deletions metrics/counter/access/daily_accumulator.py
Original file line number Diff line number Diff line change
@@ -1,33 +1,256 @@
from struct import Struct

_ACCESS_KEY = Struct("!10I")
_SESSION_KEY = Struct("!5I")
_NONE_METADATA_ID = 0
_EMPTY_DOCUMENT_ID = 1


class _CompactAccessRecord:
__slots__ = (
"access_method",
"access_date",
"access_month",
"access_year",
"collection",
"content_language",
"content_type",
"counter_access_type",
"country_code",
"document",
"document_type",
"first_second",
"first_url",
"media_format",
"multiple_timestamps",
"pid_generic",
"pid_v2",
"pid_v3",
"publication_year",
"session",
"source",
"source_key",
"title_pid_generic",
)

def __init__(self, accumulator, data, session, url, second):
self.collection = accumulator._intern(data.get("collection"))
self.source_key = accumulator._intern(data.get("source_key"))
self.document_type = accumulator._intern(data.get("document_type"))
self.pid_v2 = accumulator._intern(data.get("pid_v2"))
self.pid_v3 = accumulator._intern(data.get("pid_v3"))
self.pid_generic = accumulator._intern(data.get("pid_generic"))
self.title_pid_generic = accumulator._intern(data.get("title_pid_generic"))
self.media_format = accumulator._intern(data.get("media_format"))
self.content_language = accumulator._intern(data.get("content_language"))
self.content_type = accumulator._intern(data.get("content_type"))
self.country_code = accumulator._intern(data.get("access_country_code"))
self.access_date = accumulator._intern(data.get("access_date"))
self.access_year = accumulator._intern(data.get("access_year"))
self.access_month = accumulator._intern(data.get("access_month"))
self.publication_year = accumulator._intern(data.get("publication_year"))
self.counter_access_type = accumulator._intern(data.get("counter_access_type"))
self.access_method = accumulator._intern(data.get("access_method"))
self.source = accumulator._intern_source(
data.get("source_key"),
data.get("source"),
)
self.document = accumulator._intern_document(data, share_empty=True)
self.session = session
self.first_url = accumulator._intern(url)
self.first_second = second
self.multiple_timestamps = None

def add_timestamp(self, url, second):
if self.multiple_timestamps is None:
if url == self.first_url and second == self.first_second:
return
self.multiple_timestamps = {self.first_url: self.first_second}

current = self.multiple_timestamps.get(url)
if current is None:
self.multiple_timestamps[url] = second
elif isinstance(current, int):
if current != second:
self.multiple_timestamps[url] = {current, second}
else:
current.add(second)

def as_dict(self, accumulator):
return {
"collection": accumulator._resolve(self.collection),
"source_key": accumulator._resolve(self.source_key),
"document_type": accumulator._resolve(self.document_type),
"pid_v2": accumulator._resolve(self.pid_v2),
"pid_v3": accumulator._resolve(self.pid_v3),
"pid_generic": accumulator._resolve(self.pid_generic),
"document": accumulator._documents[self.document],
"title_pid_generic": accumulator._resolve(self.title_pid_generic),
"user_session_id": self.session,
"click_timestamps_by_url": self._timestamps_as_dict(accumulator),
"media_format": accumulator._resolve(self.media_format),
"content_language": accumulator._resolve(self.content_language),
"content_type": accumulator._resolve(self.content_type),
"access_country_code": accumulator._resolve(self.country_code),
"access_date": accumulator._resolve(self.access_date),
"access_year": accumulator._resolve(self.access_year),
"access_month": accumulator._resolve(self.access_month),
"publication_year": accumulator._resolve(self.publication_year),
"counter_access_type": accumulator._resolve(self.counter_access_type),
"access_method": accumulator._resolve(self.access_method),
"source": accumulator._sources[self.source],
}

def _timestamps_as_dict(self, accumulator):
if self.multiple_timestamps is None:
return {
accumulator._resolve(self.first_url): {self.first_second: 1},
}

timestamps = {}
for url, seconds in self.multiple_timestamps.items():
if isinstance(seconds, int):
seconds = (seconds,)
timestamps[accumulator._resolve(url)] = {
second: 1 for second in sorted(seconds)
}
return timestamps


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

def __init__(self):
super().__init__()
self._documents = {}
self._sources = {}
self._documents = [None, {}]
self._document_ids = {}
self._sources = [None]
self._source_ids = {}
self._sessions = {}
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:
value["source"] = self._sources.setdefault(source_key, source)

document_key = (
value.get("document_type"),
value.get("pid_v2"),
value.get("pid_v3"),
value.get("pid_generic"),
value.get("title_pid_generic"),
)
source_id = self._intern_source(source_key, source)
value["source"] = self._sources[source_id]

document_key = self._document_key(value)
document = value.get("document")
document_identifiers = document_key[1:]
if document is not None and any(document_identifiers):
value["document"] = self._documents.setdefault(document_key, 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._sessions.setdefault(
user_session_id,
user_session_id,
)
value["user_session_id"] = self._legacy_intern_session(user_session_id)

super().__setitem__(key, value)

def accumulate_access(self, data, session_key, url, second):
session = self._intern_session(session_key)
key = _ACCESS_KEY.pack(
self._intern(data.get("collection")),
self._intern(data.get("source_key")),
self._intern(data.get("pid_v2")),
self._intern(data.get("pid_v3")),
self._intern(data.get("pid_generic")),
session,
self._intern(data.get("access_country_code")),
self._intern(data.get("content_language")),
self._intern(data.get("media_format")),
self._intern(data.get("content_type")),
)
record = dict.get(self, key)
if record is None:
record = _CompactAccessRecord(self, data, session, url, second)
dict.__setitem__(self, 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):
yield value.as_dict(self)
else:
yield value

def _intern(self, value):
if value is None:
return 0
value_id = self._string_ids.get(value)
if value_id is None:
value_id = len(self._strings)
self._string_ids[value] = value_id
self._strings.append(value)
return value_id

def _resolve(self, value_id):
return self._strings[value_id]

def _intern_session(self, session_key):
compact_key = _SESSION_KEY.pack(
self._intern(session_key[0]),
self._intern(session_key[1]),
self._intern(session_key[2]),
session_key[3],
session_key[4],
)
session_id = self._sessions.get(compact_key)
if session_id is None:
session_id = len(self._sessions) + 1
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
if not source_key:
self._sources.append(source)
return len(self._sources) - 1
source_id = self._source_ids.get(source_key)
if source_id is None:
source_id = len(self._sources)
self._source_ids[source_key] = source_id
self._sources.append(source)
return source_id

def _intern_document(self, data, share_empty=False):
document_key = self._document_key(data)
document = data.get("document")
if document is None:
return _NONE_METADATA_ID
if share_empty and not document:
return _EMPTY_DOCUMENT_ID
if not any(document_key[1:]):
self._documents.append(document)
return len(self._documents) - 1
document_id = self._document_ids.get(document_key)
if document_id is None:
document_id = len(self._documents)
self._document_ids[document_key] = document_id
self._documents.append(document)
return document_id

@staticmethod
def _document_key(data):
return (
data.get("document_type"),
data.get("pid_v2"),
data.get("pid_v3"),
data.get("pid_generic"),
data.get("title_pid_generic"),
)
3 changes: 2 additions & 1 deletion metrics/counter/indexing/converter.py
Original file line number Diff line number Diff line change
Expand Up @@ -27,8 +27,9 @@ def convert(data):
def _convert_granularity(data, granularity):
converted_data = {}
unique_state = _initialize_unique_state()
values = getattr(data, "iter_materialized_values", data.values)

for value in data.values():
for value in values():
pipeline = _get_pipeline(value)
pipeline.accumulate(
data=converted_data,
Expand Down
Loading
Loading