Files
catalogue_backend/orchestration/partitions.py

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