71 lines
2.5 KiB
Python
71 lines
2.5 KiB
Python
"""Brand partitioning.
|
|
|
|
One Dagster partition per active brand. This is the natural expression of the
|
|
"Brand Selection -> Ingestion -> ..." flow, and it buys three concrete things
|
|
that a single un-partitioned asset would not:
|
|
|
|
* per-brand lineage in the UI, so "which brands are stale" is readable at a
|
|
glance rather than buried in one run's logs;
|
|
* per-brand retry - a failed Nestle partition does not force Amul to be
|
|
rebuilt;
|
|
* bounded memory. Each run holds one brand's rows, not the whole catalog,
|
|
which is what keeps this comfortable on an 8GB machine.
|
|
|
|
The partition set is derived from ACTIVE_BRANDS, so widening the working set
|
|
from 3 brands to 25 adds partitions with no code change.
|
|
"""
|
|
# NOTE: deliberately no `from __future__ import annotations` here.
|
|
# Dagster resolves the decorated function signatures at definition time to
|
|
# validate the `context` parameter and to infer asset input types. Under
|
|
# PEP 563/649 the annotations arrive as strings and that validation fails
|
|
# with "Cannot annotate `context` parameter with type AssetExecutionContext".
|
|
# Local Python is 3.14, which defers annotations by default, so this is not
|
|
# hypothetical.
|
|
|
|
from typing import List
|
|
|
|
from dagster import StaticPartitionsDefinition
|
|
|
|
|
|
def active_brand_names() -> List[str]:
|
|
"""Brand names to partition over.
|
|
|
|
Falls back to the brand tables that exist when ACTIVE_BRANDS is unset, and
|
|
finally to a single literal so that the definitions still LOAD when the
|
|
database is unreachable. A Dagster code location that cannot be loaded is
|
|
far harder to debug than one that loads and shows an empty run - and the
|
|
webserver imports this at startup, before anyone can fix a connection.
|
|
"""
|
|
from app.services.active_brands import active_display_names
|
|
|
|
names = active_display_names()
|
|
if names:
|
|
return names
|
|
|
|
try:
|
|
from app.services.vector_store import list_available_brands
|
|
|
|
live = list_available_brands()
|
|
if live:
|
|
return live
|
|
except Exception: # noqa: BLE001 - see docstring
|
|
pass
|
|
|
|
return ["Amul"]
|
|
|
|
|
|
_PARTITIONS = None
|
|
|
|
|
|
def brand_partitions() -> StaticPartitionsDefinition:
|
|
"""Cached so every asset shares one identical partition set.
|
|
|
|
Dagster compares partition definitions by value, but building the list once
|
|
also means a single database round trip at code-load time instead of one
|
|
per asset.
|
|
"""
|
|
global _PARTITIONS
|
|
if _PARTITIONS is None:
|
|
_PARTITIONS = StaticPartitionsDefinition(active_brand_names())
|
|
return _PARTITIONS
|