Files
catalogue_backend/orchestration/jobs.py
2026-08-28 07:49:56 +05:30

107 lines
4.0 KiB
Python

"""Jobs - four logical workflows, deliberately not one large one.
Splitting them this way follows how the work actually differs in cost and
cadence:
* ingestion is cheap and offline (seed files), so it can run often;
* embedding drives the sentence-transformer, so it is separated to be
retried or run alone without redoing ingestion;
* nutrition enrichment makes third-party calls and is the slowest thing here;
* ML training needs order history, which is unrelated to catalog freshness.
Selecting a subset of assets in one giant job would express the same graph,
but you could not schedule the parts on different cadences, and a failure in
one concern would show as a failure of everything.
"""
# 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 dagster import AssetSelection, define_asset_job
# NOTE: no partitions_def on define_asset_job. Dagster infers the partitioning
# from the selected assets, and passing it explicitly is deprecated (removed in
# 2.0). The catalog and embedding jobs are still per-brand partitioned because
# every asset they select is.
# Brand -> raw -> validated -> enriched -> Postgres.
catalog_ingestion_job = define_asset_job(
name="catalog_ingestion_job",
description=(
"Brand selection, seed intake, validation, enrichment and the "
"Postgres upsert, for one brand partition."
),
selection=AssetSelection.assets(
"active_brand",
"raw_products",
"validated_products",
"enriched_products",
"catalog_database",
),
)
# Stored rows -> embeddings -> a verified vector index.
embedding_refresh_job = define_asset_job(
name="embedding_refresh_job",
description=(
"Generate embeddings for stored rows that lack one, then verify the "
"brand is fully searchable in pgvector."
),
selection=AssetSelection.assets("product_embeddings", "vector_index"),
)
# Stored rows -> Open Food Facts -> nutrition_facts/insights -> the 2 models.
nutrition_enrichment_job = define_asset_job(
name="nutrition_enrichment_job",
description=(
"Fetch verified nutrition for a brand's products, score them, then "
"refit the similarity and clustering models across all brands."
),
selection=AssetSelection.assets("nutrition_data", "nutrition_models"),
)
# Synthetic order history -> fitted models -> artifact report.
ml_training_job = define_asset_job(
name="ml_training_job",
description=(
"Rebuild the store-intelligence training data, fit the served models "
"(discount, trending, popularity) and report the artifacts on disk."
),
selection=AssetSelection.assets(
"training_dataset", "trained_models", "model_evaluation"
),
)
# Uploaded spreadsheets -> the same 11 stages -> brand tables.
#
# Separate from catalog_ingestion_job because the source is different, not
# because the work is: that job reads seed JSON for a brand, this one reads a
# batch of files somebody uploaded. They share every stage in between, and
# keeping them apart means a failure here does not read as the brand pipeline
# being broken.
batch_ingestion_job = define_asset_job(
name="batch_ingestion_job",
description=(
"Ingest a batch of uploaded store spreadsheets: resolve the batch, "
"parse each file, run the 11 stages, and report what landed."
),
selection=AssetSelection.assets(
"batch_manifest",
"batch_parsed_files",
"batch_catalog_rows",
"batch_ingest_report",
),
)
ALL_JOBS = [
catalog_ingestion_job,
embedding_refresh_job,
nutrition_enrichment_job,
ml_training_job,
batch_ingestion_job,
]