"""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, ]