Backend- file ingestion API Updates

This commit is contained in:
sriram
2026-08-28 07:49:56 +05:30
parent f698720ee2
commit a54bd43f8b
17 changed files with 2306 additions and 138 deletions

View File

@@ -144,10 +144,21 @@ If you find yourself writing business logic in this package, it belongs in
| `embedding_refresh_job` | embed rows lacking a vector, verify the index |
| `nutrition_enrichment_job` | fetch nutrition, score it, refit the two models |
| `ml_training_job` | rebuild the store dataset, fit the served models, report |
| `batch_ingestion_job` | a batch of uploaded spreadsheets -> the 11 stages -> Postgres |
Four jobs rather than one because they differ in cost and cadence: ingestion is
Five jobs rather than one because they differ in cost and cadence: ingestion is
offline and cheap, embedding drives the transformer, nutrition leaves the
machine, and ML training depends on order history rather than catalog freshness.
machine, ML training depends on order history rather than catalog freshness, and
batch ingestion is event-driven - it exists only when somebody uploads files.
`batch_ingestion_job` shares every stage with the upload path in production;
both call `app/core/batch_ingest.py`. Pick a batch with run config:
```json
{"ops": {"batch_manifest": {"config": {"batch_id": "<id from the Admin UI>"}}}}
```
Left empty it takes the oldest batch still waiting under `BATCH_UPLOAD_DIR`.
From the command line:
@@ -170,6 +181,7 @@ dagster asset materialize -m orchestration.definitions \
| `daily_nutrition_refresh` | 04:00 | slowest, and the only one making outbound calls |
| `weekly_ml_retrain` | Sun 05:00 | training data is a deterministic simulation, so nightly would repeat itself |
| `seed_catalog_sensor` | every 60s | re-ingest a brand when its seed file changes |
| `batch_upload_sensor` | every 60s | run an uploaded batch that is staged and waiting |
Shipping them stopped is deliberate. This machine also runs Postgres, Ollama,
the API and Vite; a schedule that started itself the moment `dagster dev`
@@ -239,6 +251,17 @@ allocates the whole 8GB VPS (ollama 3G, backend 2560M, postgres 1G), so there
is no headroom for a webserver plus daemon. Orchestration is a development
concern here.
**This includes `batch_ingestion_job`.** The Batch Catalog Ingestion screen on
the site does not talk to Dagster and does not need it running. It stages the
uploaded files, then executes the same functions this job wraps
(`app/core/batch_ingest.py`) on a single bounded worker thread inside the API
process - see `app/core/batch_worker.py` for why one worker rather than a
thread per upload. The job here exists so the graph is inspectable and a batch
can be re-run from the Launchpad on a development machine.
Do not turn `batch_upload_sensor` on next to a running API: both would pick up
the same staged batch. That is why it ships STOPPED like the rest.
Port 3030, not Dagster's default 3000 - `serve.py` binds `PORTS=3000,8000` and
the frontend nginx also listens on 3000.

View File

@@ -1,3 +1,3 @@
from orchestration.assets import catalog, ml, nutrition
from orchestration.assets import batch_catalog, catalog, ml, nutrition
__all__ = ["catalog", "nutrition", "ml"]
__all__ = ["catalog", "nutrition", "ml", "batch_catalog"]

View File

@@ -0,0 +1,315 @@
"""Batch ingestion assets: a staged upload -> parsed -> ingested -> reported.
Every asset delegates to `app.core.batch_ingest`, which is the same module the
production endpoint calls. What Dagster adds here is the lineage between the
steps, per-run metadata that survives a restart, retries on the step that
reaches the network, and a launchpad for re-running a batch that went wrong -
none of which the worker thread and its in-memory store can give.
WHY THESE ASSETS ARE NOT PARTITIONED
------------------------------------
The brand assets in `catalog.py` are partitioned per brand, because a brand is
a stable, enumerable thing and one partition per brand buys per-brand retry and
bounded memory. A batch is neither stable nor enumerable in advance - it comes
into existence when somebody uploads files. Modelling it as a partition would
mean adding a dynamic partition per upload and never removing it, so the
partition set grows forever and the UI fills with dead keys. The batch id lives
in run config instead, which is what run config is for.
WHERE THIS RUNS
---------------
`dagster dev` on a developer's machine. Dagster is not deployed: it is absent
from docker-compose.prod.yml, its dependencies are in a separate requirements
file the API image does not install, and it is not on any request path. The
production path executes the identical functions on a bounded worker thread
(app/core/batch_worker.py). That is deliberate - a webserver plus daemon costs
400-600MB resident, and the host this ships to does not have it to give.
"""
# 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.
import time
from typing import Any, Dict, List
from dagster import (
AssetExecutionContext,
Backoff,
Failure,
Jitter,
MetadataValue,
RetryPolicy,
asset,
)
from orchestration.config import BatchConfig, database_target, require_local_database
# Same policy and same reasoning as catalog.py: only the step that leaves the
# machine is retried. A file that fails to parse will fail identically on a
# second attempt, and retrying it only delays the red run somebody has to read.
NETWORK_RETRY = RetryPolicy(
max_retries=2, delay=5, backoff=Backoff.EXPONENTIAL, jitter=Jitter.PLUS_MINUS
)
GROUP = "batch_catalog"
def _pick_batch_id(config: BatchConfig) -> str:
"""Run config wins; otherwise the oldest batch still waiting to run."""
from app.core import batch_ingest
if config.batch_id:
return config.batch_id
waiting = [
m for m in batch_ingest.list_manifests()
if m.status in (batch_ingest.QUEUED, batch_ingest.INTERRUPTED)
]
if not waiting:
raise Failure(
description=(
"No batch_id given and no staged batch is waiting to run.\n\n"
"Upload one through Admin -> Batch Catalog Ingestion on the site, "
"or pass {\"batch_id\": \"...\"} in the Launchpad. Staged batches "
"live under BATCH_UPLOAD_DIR ({}).".format(batch_ingest.batch_root())
),
metadata={"staged_batches": len(batch_ingest.list_manifests())},
)
return sorted(waiting, key=lambda m: m.created_at)[0].batch_id
@asset(
group_name=GROUP,
description="The staged upload batch this run will ingest.",
)
def batch_manifest(context: AssetExecutionContext, config: BatchConfig) -> Dict[str, Any]:
"""Root of the lineage: resolve which batch, and prove it is readable.
Resolving here rather than inside each downstream asset means the run page
shows which files a run actually touched, up front, instead of it being
discovered halfway through.
"""
from app.core import batch_ingest
batch_id = _pick_batch_id(config)
manifest = batch_ingest.read_manifest(batch_id)
if manifest is None:
raise Failure(
description="No manifest for batch '{}' under {}.".format(
batch_id, batch_ingest.batch_root()
),
metadata={"batch_id": batch_id},
)
pending = [f for f in manifest.files if f.status == batch_ingest.QUEUED]
context.add_output_metadata(
{
"batch_id": batch_id,
"status": manifest.status,
"files_total": manifest.files_total,
"files_pending": len(pending),
"use_llm": config.use_llm or manifest.use_llm,
"fetch_images": config.fetch_images or manifest.fetch_images,
"database": database_target(),
"files": MetadataValue.md(
"\n".join(
"- `{}` - {}".format(f.filename, f.status) for f in manifest.files
) or "_none_"
),
}
)
return {"batch_id": batch_id, "files_pending": len(pending)}
@asset(
group_name=GROUP,
description="Column mapping and row count for each staged file.",
)
def batch_parsed_files(
context: AssetExecutionContext, batch_manifest: Dict[str, Any]
) -> List[Dict[str, Any]]:
"""Parse every staged file without ingesting anything.
This is the Dagster equivalent of the UI's Preview step, and it exists for
the same reason: the mapping from a store's headers onto catalog fields is
a guess, and it is much cheaper to see it now than after eleven stages have
run over two thousand rows.
A file that will not parse is reported, not raised. One unreadable sheet
among ten is a data problem worth seeing, not a reason to block the nine.
"""
from app.core import batch_ingest
from app.core import store_catalog_pipeline as pipeline
batch_id = batch_manifest["batch_id"]
manifest = batch_ingest.read_manifest(batch_id)
directory = batch_ingest.batch_dir(batch_id)
parsed: List[Dict[str, Any]] = []
lines = []
for entry in manifest.files:
if not entry.stored_name:
parsed.append({"filename": entry.filename, "ok": False,
"error": entry.detail or "not staged"})
lines.append("- `{}` - **not staged**".format(entry.filename))
continue
try:
content = (directory / entry.stored_name).read_bytes()
frame, mapping = pipeline.parse_spreadsheet(entry.filename, content)
except Exception as exc: # noqa: BLE001 - see docstring
parsed.append({"filename": entry.filename, "ok": False, "error": str(exc)})
lines.append("- `{}` - **unreadable**: {}".format(entry.filename, exc))
continue
recognised = {f: str(c) for f, c in mapping.columns.items()}
parsed.append({
"filename": entry.filename,
"ok": True,
"rows": int(len(frame)),
"recognised_columns": recognised,
"unrecognised_columns": list(mapping.unrecognised),
})
lines.append("- `{}` - {} rows, {} mapped column(s)".format(
entry.filename, len(frame), len(recognised)
))
context.add_output_metadata(
{
"batch_id": batch_id,
"files": len(parsed),
"files_readable": sum(1 for p in parsed if p["ok"]),
"rows_total": sum(int(p.get("rows") or 0) for p in parsed),
"detail": MetadataValue.md("\n".join(lines) or "_none_"),
}
)
return parsed
@asset(
group_name=GROUP,
retry_policy=NETWORK_RETRY,
description="The 11-stage pipeline run over every pending file in the batch.",
)
def batch_catalog_rows(
context: AssetExecutionContext,
config: BatchConfig,
batch_manifest: Dict[str, Any],
batch_parsed_files: List[Dict[str, Any]],
) -> Dict[str, Any]:
"""Stages 1-11, per file, via the shared `batch_ingest.run_batch`.
THE WRITE GUARD IS CALLED HERE, not only in the reporting asset downstream.
`run_batch` reaches stage 11 and upserts, so by the time the next asset runs
the write has already happened - checking there would be checking after the
fact. backend/.env points at production; see orchestration/config.py.
Per-file progress is logged rather than pushed anywhere. In production the
same callback feeds the polling endpoint; in Dagster the run log IS the
progress view, and writing to both would be two sources of truth.
"""
from app.core import batch_ingest
target = require_local_database("batch_catalog_rows")
batch_id = batch_manifest["batch_id"]
started = time.time()
# Run config overrides what the upload asked for, so a batch uploaded with
# images off can be re-run with them on without re-uploading.
manifest = batch_ingest.read_manifest(batch_id)
manifest.use_llm = bool(config.use_llm or manifest.use_llm)
manifest.fetch_images = bool(config.fetch_images or manifest.fetch_images)
batch_ingest.write_manifest(manifest)
seen = {"file": None}
def on_change(current):
name = current.current_file
if name and name != seen["file"]:
seen["file"] = name
context.log.info("Ingesting %s", name)
result = batch_ingest.run_batch(batch_id, on_change=on_change)
totals = result.totals()
context.add_output_metadata(
{
"batch_id": batch_id,
"status": result.status,
"files_total": result.files_total,
"files_done": result.files_done,
"files_failed": result.files_failed,
"rows_ingested": totals["rows_total"],
"products_built": totals["products_built"],
"inserted": totals["inserted"],
"backfilled": totals["backfilled"],
"unchanged": totals["skipped_existing"],
"rejected": totals["rejected"],
"brands": ", ".join(result.brands()) or "none",
"database": target,
"cleanup": "False (never deletes rows outside this batch)",
"duration_s": round(time.time() - started, 2),
}
)
return result.to_dict()
@asset(
group_name=GROUP,
description="Per-file outcome of the batch, and a red run if none of it landed.",
)
def batch_ingest_report(
context: AssetExecutionContext, batch_catalog_rows: Dict[str, Any]
) -> Dict[str, Any]:
"""Turn the batch result into a readable report, and decide run colour.
A batch where every file failed materializes GREEN without this: nothing
raised, so from Dagster's point of view the work completed. That is the
worst outcome available here - the run looks fine and the catalog did not
change - so it is made loud, in the same spirit as `raw_products` failing on
an empty brand.
A partial batch stays green. Four files landing out of five is a data
problem to read about, not a pipeline failure, and failing the run would
make the four look like they had not happened.
"""
files = batch_catalog_rows.get("files") or []
done = [f for f in files if f.get("status") == "done"]
failed = [f for f in files if f.get("status") == "failed"]
lines = []
for entry in files:
lines.append("- `{}` - **{}** - {}".format(
entry.get("filename"), entry.get("status"), entry.get("detail") or ""
))
context.add_output_metadata(
{
"batch_id": batch_catalog_rows.get("batch_id"),
"status": batch_catalog_rows.get("status"),
"files_done": len(done),
"files_failed": len(failed),
"totals": MetadataValue.json(batch_catalog_rows.get("totals") or {}),
"report": MetadataValue.md("\n".join(lines) or "_no files_"),
}
)
if files and not done:
raise Failure(
description=(
"Every file in batch '{}' failed; nothing reached the catalog.".format(
batch_catalog_rows.get("batch_id")
)
),
metadata={"files_failed": len(failed)},
)
return {
"batch_id": batch_catalog_rows.get("batch_id"),
"status": batch_catalog_rows.get("status"),
"files_done": len(done),
"files_failed": len(failed),
}

View File

@@ -94,6 +94,25 @@ class BrandConfig(Config):
max_products_per_brand: Optional[int] = None
class BatchConfig(Config):
"""Per-run config for the uploaded-spreadsheet batch job.
`batch_id` selects a batch staged under BATCH_UPLOAD_DIR by the API. Left
empty, the job picks the oldest batch still queued, which is what the
sensor wants and what makes a manual "run it now" from the Launchpad
convenient.
The two network flags default OFF for the same reason they do on
BrandConfig, with more force: a batch is up to twenty files, so leaving
image search on would fire tens of thousands of outbound requests from one
run.
"""
batch_id: Optional[str] = None
use_llm: bool = False
fetch_images: bool = False
def resolve_brands(cfg_brands: Optional[List[str]]) -> List[str]:
"""Run config wins; otherwise ACTIVE_BRANDS; otherwise whatever the DB has.

View File

@@ -61,11 +61,11 @@ _ENV_SOURCE = _load_orchestration_env()
from dagster import Definitions, load_assets_from_modules, multiprocess_executor # noqa: E402
from orchestration import config as orchestration_config # noqa: E402
from orchestration.assets import catalog, ml, nutrition # noqa: E402
from orchestration.assets import batch_catalog, catalog, ml, nutrition # noqa: E402
from orchestration.jobs import ALL_JOBS # noqa: E402
from orchestration.schedules import ALL_SCHEDULES, ALL_SENSORS # noqa: E402
_all_assets = load_assets_from_modules([catalog, nutrition, ml])
_all_assets = load_assets_from_modules([catalog, nutrition, ml, batch_catalog])
defs = Definitions(
assets=_all_assets,

View File

@@ -76,9 +76,31 @@ ml_training_job = define_asset_job(
),
)
# 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,
]

View File

@@ -31,6 +31,7 @@ from dagster import (
)
from orchestration.jobs import (
batch_ingestion_job,
catalog_ingestion_job,
embedding_refresh_job,
ml_training_job,
@@ -149,4 +150,66 @@ def seed_catalog_sensor(context: SensorEvaluationContext):
]
ALL_SENSORS = [seed_catalog_sensor]
@sensor(
job=batch_ingestion_job,
minimum_interval_seconds=60,
default_status=DefaultSensorStatus.STOPPED,
description="Run a batch of uploaded spreadsheets once one is staged and waiting.",
)
def batch_upload_sensor(context: SensorEvaluationContext):
"""Watch BATCH_UPLOAD_DIR for batches the API staged but did not run.
There is NO schedule for batch ingestion, deliberately. A batch exists
because a person uploaded files; there is nothing to do on a cron, and a
schedule that woke up hourly to find nothing would be pure cost on a
machine that is already sharing itself with Postgres, the API and Vite.
Like seed_catalog_sensor this is as simple as it can be - a directory
listing once a minute, no watcher process, no broker. The cursor is the set
of batch ids already requested, so a batch is launched once even though it
stays on disk afterwards.
NOTE ON THE NORMAL PRODUCTION PATH: nothing here is involved. The API runs
a staged batch itself, on a bounded worker thread, and this sensor exists
for the development machine where Dagster is the thing driving the work.
Turning it on alongside a running API would mean both trying to ingest the
same batch, which is why it ships STOPPED.
"""
from app.core import batch_ingest
root = batch_ingest.batch_root()
if not root.exists():
return SkipReason("Batch upload directory {} does not exist.".format(root))
waiting = [
m for m in batch_ingest.list_manifests()
if m.status in (batch_ingest.QUEUED, batch_ingest.INTERRUPTED)
]
if not waiting:
return SkipReason("No staged batch is waiting to run.")
already = set(filter(None, (context.cursor or "").split(",")))
fresh = [m for m in waiting if m.batch_id not in already]
if not fresh:
return SkipReason(
"{} staged batch(es), all already requested.".format(len(waiting))
)
# Keep the cursor bounded - it is a string in the Dagster instance, not a
# log, and only the recent tail is ever consulted.
context.update_cursor(",".join(list(already | {m.batch_id for m in fresh})[-200:]))
context.log.info(
"Requesting runs for staged batch(es): %s",
", ".join(m.batch_id for m in fresh),
)
return [
RunRequest(
run_key=m.batch_id,
run_config={"ops": {"batch_manifest": {"config": {"batch_id": m.batch_id}}}},
)
for m in fresh
]
ALL_SENSORS = [seed_catalog_sensor, batch_upload_sensor]