Dagster Orchestration

This commit is contained in:
sriram
2026-08-29 14:49:45 +05:30
parent 27d53fa957
commit 998df898db
28 changed files with 1455 additions and 76 deletions

View File

@@ -28,7 +28,7 @@ from dataclasses import dataclass
from typing import List, Optional, Tuple
from fastapi import HTTPException, UploadFile
from pydantic import BaseModel
from pydantic import BaseModel, Field
from app.api.batch_job_store import batch_job_store
from app.core import batch_ingest, batch_worker
@@ -182,6 +182,22 @@ def parse_all(read: List[Tuple[str, bytes]], limits: UploadLimits):
# ---------------------------------------------------------------------------
# Response shape
# ---------------------------------------------------------------------------
class StageOut(BaseModel):
"""One pipeline stage as it happened to one file.
Kept even in slim list responses: eleven of these per file is a few hundred
bytes, unlike the `products` manifest slim exists to drop.
"""
index: int
name: str
rows_done: int = 0
rows_total: int = 0
started_at: Optional[float] = None
# None while the stage is still running.
finished_at: Optional[float] = None
class BatchFileOut(BaseModel):
index: int
filename: str
@@ -196,7 +212,12 @@ class BatchFileOut(BaseModel):
total_stages: int = pipeline.TOTAL_STAGES
rows_done: int = 0
rows_total: int = 0
# The stages this file has been through, so a FINISHED file can still show
# its timeline. The scalars above only ever say where it is right now.
stages: List[StageOut] = Field(default_factory=list)
size_bytes: int = 0
started_at: Optional[float] = None
finished_at: Optional[float] = None
result: Optional[dict] = None
@@ -213,6 +234,14 @@ class BatchOut(BaseModel):
current_file: Optional[str] = None
use_llm: bool
fetch_images: bool
# Which executor owns this batch - "inprocess" or "dagster". A dagster batch
# sits queued until the orchestrator picks it up, which the UI has to be
# able to say out loud rather than showing a run that looks stuck.
runner: str = batch_ingest.RUNNER_INPROCESS
# The 11 stage names, in order, so a client can draw the whole pipeline
# before a file has entered any of it. Served rather than duplicated in the
# frontend so the two cannot drift when a stage is added.
stage_names: List[str] = Field(default_factory=lambda: list(pipeline.STAGE_NAMES))
totals: dict
brands: List[str]
files: List[BatchFileOut]
@@ -287,6 +316,49 @@ def stage_and_queue(
return manifest, True
def stage_for_orchestrator(
valid: List[Tuple[str, bytes, int]],
invalid: List[Tuple[str, str]],
*,
use_llm: bool,
fetch_images: bool,
submitted_by: Optional[str] = None,
) -> batch_ingest.BatchManifest:
"""Publish a batch for Dagster to claim, and hand it to no one else.
Same staging as `stage_and_queue`, minus the `batch_worker.submit`. The
batch is left `queued` and stamped `runner="dagster"`, which is what
`_pick_batch_id` and `batch_upload_sensor` filter on; the in-process worker
never scans for work, so leaving it unsubmitted is enough to keep the two
executors off each other's batches.
A third function rather than a `runner=` argument on `stage_and_queue`, for
the reason given in `stage_pending`: whether work starts here or somewhere
else is not the kind of decision that should hang off a boolean anyone can
flip later.
The batch sits queued until an orchestrator actually runs - which, if
`dagster dev` is not up, is never. Callers are expected to say so rather
than present it as a run in progress.
"""
manifest = batch_ingest.stage_batch(
[(name, contents) for name, contents, _n in valid],
use_llm=use_llm,
fetch_images=fetch_images,
invalid=invalid,
submitted_by=submitted_by,
)
for entry, (_name, _contents, rows) in zip(manifest.files, valid):
entry.rows_total = rows
manifest.runner = batch_ingest.RUNNER_DAGSTER
manifest.detail = (
"Waiting for the Dagster orchestrator to pick this batch up."
)
batch_ingest.write_manifest(manifest)
batch_job_store.put(manifest)
return manifest
def stage_pending(
valid: List[Tuple[str, bytes, int]],
invalid: List[Tuple[str, str]],

View File

@@ -202,7 +202,13 @@ def get_catalog_batch(batch_id: str) -> BatchOut:
@router.post("/batches/{batch_id}/resume", dependencies=[Depends(require_admin)])
def resume_catalog_batch(batch_id: str) -> BatchOut:
"""Re-queue a batch a restart cut short, or one that was queued behind a full queue."""
"""Re-queue a batch a restart cut short, or one that was queued behind a full queue.
Also the way out of a batch staged for Dagster that no orchestrator ever
came for - the "run it here instead" button. Because this hands the batch to
THIS container's worker, it also takes ownership: the runner is flipped to
`inprocess` so Dagster will not claim a batch that is already running here.
"""
manifest = batch_ingest.read_manifest(batch_id)
if not manifest:
raise HTTPException(status_code=404, detail="Batch not found")
@@ -217,6 +223,7 @@ def resume_catalog_batch(batch_id: str) -> BatchOut:
batch_job_store.clear_cancel(batch_id)
manifest.status = batch_ingest.QUEUED
manifest.detail = None
manifest.runner = batch_ingest.RUNNER_INPROCESS
batch_ingest.write_manifest(manifest)
batch_job_store.put(manifest)
@@ -314,6 +321,10 @@ class InboxStartRequest(InboxSelection):
# belongs to the person who can see what the machine is already doing.
use_llm: bool = False
fetch_images: bool = False
# Who runs it. "inprocess" is this container's worker thread and is the
# default, so an existing client that never sends the field is unaffected.
# "dagster" stages the batch and leaves it for the orchestrator to claim.
runner: str = batch_ingest.RUNNER_INPROCESS
class InboxDismissOut(BaseModel):
@@ -405,6 +416,14 @@ def start_batch_from_inbox(request: InboxStartRequest) -> BatchOut:
The originals are removed afterwards, so the same sheet cannot be started
twice from a stale checkbox in another tab.
"""
if request.runner not in batch_ingest.RUNNERS:
raise HTTPException(
status_code=400,
detail="Unknown runner {!r}. Expected one of: {}.".format(
request.runner, ", ".join(sorted(batch_ingest.RUNNERS))
),
)
grouped = _parse_file_ids(request.file_ids)
if not grouped:
raise HTTPException(status_code=400, detail="No files were selected.")
@@ -447,23 +466,35 @@ def start_batch_from_inbox(request: InboxStartRequest) -> BatchOut:
unique_senders = sorted(set(senders))
submitted_by = ", ".join(unique_senders)[:120] if unique_senders else None
manifest, started = batch_common.stage_and_queue(
picked,
[],
use_llm=request.use_llm,
fetch_images=request.fetch_images,
submitted_by=submitted_by,
)
if not started:
raise HTTPException(
status_code=429,
detail=(
"Too many batches are already queued. These files have been "
"taken out of the inbox and saved as batch "
f"{manifest.batch_id} - press Resume on it once the current "
"batch finishes."
),
if request.runner == batch_ingest.RUNNER_DAGSTER:
# Staged and left alone: Dagster claims it on its next sensor tick, or
# from the Launchpad. Nothing here waits on that, and the batch is
# durable either way.
manifest = batch_common.stage_for_orchestrator(
picked,
[],
use_llm=request.use_llm,
fetch_images=request.fetch_images,
submitted_by=submitted_by,
)
else:
manifest, started = batch_common.stage_and_queue(
picked,
[],
use_llm=request.use_llm,
fetch_images=request.fetch_images,
submitted_by=submitted_by,
)
if not started:
raise HTTPException(
status_code=429,
detail=(
"Too many batches are already queued. These files have been "
"taken out of the inbox and saved as batch "
f"{manifest.batch_id} - press Resume on it once the current "
"batch finishes."
),
)
# Only now, once the bytes are safely staged under a new id. Retiring them
# first would lose the files outright if staging then failed.

View File

@@ -173,7 +173,13 @@ _KEYWORD_RULES: Tuple[Tuple[str, Tuple[str, ...]], ...] = (
("description", ("description", "desc", "detail")),
("category", ("category", "segment")),
("brand", ("brand", "manufacturer", "company")),
("size_variants", ("size", "pack", "weight", "volume", "net qty", "quantity")),
# "net qty" / "net quantity" is the Indian labelling term for a pack size.
# A bare "Quantity" column is not: in a store sheet it is how many units the
# shop has or is ordering, and mapping it here made a case-pack count of 72
# into the pack size, which then got appended to the product name. Worse,
# mapping is first-wins by column position, so a leading "Quantity" column
# also shut out the sheet's real "Pack Size" column.
("size_variants", ("size", "pack", "weight", "volume", "net qty", "net quantity")),
("providers", ("provider", "platform", "marketplace", "available at")),
("highlights", ("highlight", "feature", "benefit")),
("nutrients", ("nutrient", "nutrition")),
@@ -186,12 +192,38 @@ _KEYWORD_RULES: Tuple[Tuple[str, Tuple[str, ...]], ...] = (
_BLANK_VALUES = frozenset({"", "nan", "none", "null", "na", "n/a", "-", "--", "#n/a"})
# Fields whose value is a NAME, and which therefore must not be fed from a
# sheet's id/code column for the same concept. The keyword rules match on
# substring, so "categoryid" satisfies the "category" rule and a column of
# 1/2/3 ends up stored as the product's category. Only these two fields need
# the guard: "hsn code" and "barcode" ARE identifier fields and must keep
# matching their rules.
_NAME_ONLY_FIELDS = frozenset({"category", "brand"})
_IDENTIFIER_SUFFIXES = ("id", "ids", "code", "codes", "no", "num", "number")
def _is_identifier_header(normalized: str) -> bool:
"""True when a header names an id/code column rather than a name column,
covering both "categoryid" and "category id" (the normalizer strips the
underscore in "category_id" to a space, and nothing at all from
"categoryid")."""
return any(
normalized.endswith(suffix) and normalized[: -len(suffix)].strip()
for suffix in _IDENTIFIER_SUFFIXES
)
def _canonical_field(normalized: str) -> Optional[str]:
exact = _EXACT_HEADERS.get(normalized)
if exact:
return exact
for canonical, keywords in _KEYWORD_RULES:
if any(keyword in normalized for keyword in keywords):
if canonical in _NAME_ONLY_FIELDS and _is_identifier_header(normalized):
# An id column for a name field: claim nothing, so the real
# name column (if the sheet has one) is still free to match and
# the id is reported back under `unrecognised` in the preview.
return None
return canonical
return None
@@ -453,14 +485,24 @@ def _build_product_dict(req: AddProductRequest, brand_parent: str,
if s3_urls:
final_image_urls = list(s3_urls)
# 2. Inherit from brand sample
if not final_image_urls and sample_existing.get("image_urls"):
final_image_urls = list(sample_existing.get("image_urls"))
# 3. Canonical S3 fallback URL
# A product's images are NOT inheritable from its brand.
#
# This used to fall back to `sample_existing["image_urls"]` - the
# `_brand_sample()` row, i.e. one arbitrary product of the brand,
# resolved once and reused for every row in the upload. It copied that
# product's photographs verbatim onto every image-less sibling, which is
# why Marie Gold and Milk Bikis both shipped carrying four
# `britannia_..._good_day_cashew_cookies_200g/` URLs while their own
# image_ids were perfectly correct. Another product's photo is never a
# defensible default for this one, so there is no fallback here.
#
# Nor is a URL invented. The old canonical fallback guessed
# `https://nearledaily.s3.ap-south-1.amazonaws.com/...` while the
# configured bucket is DigitalOcean Spaces (see S3_ENDPOINT), so it
# produced a guaranteed 404 that merely looked like an image. Leaving
# the list empty lets ProductCard render its real "no image" state.
if not final_image_urls:
canonical_s3 = f"https://nearledaily.s3.ap-south-1.amazonaws.com/daily/brands/{brand_slug}/{image_id}/image_000.jpg"
final_image_urls = [canonical_s3]
logger.info("No image found for '%s' (%s)", product_name, image_id)
primary_image_url = final_image_urls[0] if final_image_urls else None
search_text = f"{brand_parent} {product_name} {category} {description} {price_range}"