aboutsummaryrefslogtreecommitdiff
path: root/generalresearch/incite/mergers/ym_wall_summary.py
diff options
context:
space:
mode:
Diffstat (limited to 'generalresearch/incite/mergers/ym_wall_summary.py')
-rw-r--r--generalresearch/incite/mergers/ym_wall_summary.py20
1 files changed, 12 insertions, 8 deletions
diff --git a/generalresearch/incite/mergers/ym_wall_summary.py b/generalresearch/incite/mergers/ym_wall_summary.py
index 2f5995f..a01443b 100644
--- a/generalresearch/incite/mergers/ym_wall_summary.py
+++ b/generalresearch/incite/mergers/ym_wall_summary.py
@@ -1,7 +1,7 @@
from __future__ import annotations
from datetime import datetime, time, timedelta
-from typing import Literal, Type
+from typing import Literal
import dask.dataframe as dd
import pandas as pd
@@ -12,7 +12,8 @@ from generalresearch.incite.collections.thl_web import (
SessionDFCollection,
WallDFCollection,
)
-from generalresearch.incite.mergers import (
+from generalresearch.incite.exceptions import FetchError
+from generalresearch.incite.mergers.base import (
MergeCollection,
MergeCollectionItem,
MergeType,
@@ -41,6 +42,8 @@ class YMWallSummaryMergeItem(MergeCollectionItem):
ddf = wall_collection.ddf(
items=wall_items, force_rr_latest=False, include_partial=True
)
+ assert isinstance(ddf, pd.DataFrame)
+
ddf = ddf[ddf["started"].between(start, end)]
# Then we need the sessions for these wall events. They'll have started
@@ -82,18 +85,20 @@ class YMWallSummaryMergeItem(MergeCollectionItem):
class YMWallSummaryMerge(MergeCollection):
merge_type: Literal[MergeType.YM_WALL_SUMMARY] = MergeType.YM_WALL_SUMMARY
_schema = YMWallSummarySchema
- collection_item_class: Type[YMWallSummaryMergeItem] = YMWallSummaryMergeItem
+ collection_item_class: type[YMWallSummaryMergeItem] = YMWallSummaryMergeItem
items: list[YMWallSummaryMergeItem] = Field(default_factory=list)
@field_validator("offset")
def check_offset_ym_wall_summary(cls, v: str | None):
# the offset MUST be on a whole day, no hourly
+ assert v
assert v.endswith("D"), "offset must be in days"
return v
@field_validator("start")
def check_start_ym_wall_summary(cls, v: datetime | None):
# the start MUST be start on midnight exactly
+ assert v
assert v.time() == time(0, 0, 0, 0), "start must no have a time component"
return v
@@ -117,9 +122,8 @@ class YMWallSummaryMerge(MergeCollection):
# item every time build is run even if it isn't closed
# if item.should_archive():
item.fetch(wall_collection, session_collection, user_id_product)
- except (Exception,) as e:
+ except FetchError as e:
capture_exception(e)
- pass
@staticmethod
def build_groupbys(df: pd.DataFrame) -> pd.DataFrame:
@@ -177,10 +181,10 @@ class YMWallSummaryMerge(MergeCollection):
# df.to_parquet(str(self.archive_path) + ".all.parquet")
pass
- def get_counts(self, product_id):
+ def get_counts(self, product_id: str):
# examples...
product_id = ""
- df = dd.read_parquet(
+ _ = dd.read_parquet(
str(self.archive_path) + ".all.parquet",
filters=[
("product_id", "=", product_id),
@@ -188,7 +192,7 @@ class YMWallSummaryMerge(MergeCollection):
],
).compute()
country_iso = "de"
- df = dd.read_parquet(
+ _ = dd.read_parquet(
str(self.archive_path) + ".all.parquet",
filters=[
("product_id", "=", product_id),