M4: ingestion management (incremental, GS1 supplement, dedup/conflict, quality, scheduler)
- OFF incremental fetch via search API + persistent watermark (ingest_state, migration 0004) - GS1 barcode supplement adapter (offline mapping + GS1-style API) filling only gaps with field-level provenance - Non-GTIN dedup with canonical selection + merge_log; field-level conflict resolution (source trust > recency) - Quality scoring (0.4 completeness + 0.3 source trust + 0.2 multi-source + 0.1 freshness) wired into load/merge - Jobs: update_off, dedup, schedule; docs/ingestion-management.md - 19 new tests (pure + DB-integration), ruff clean Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>
This commit is contained in:
@@ -0,0 +1,40 @@
|
||||
"""Deduplicate products: merge non-GTIN duplicates into a canonical record.
|
||||
|
||||
Usage:
|
||||
python -m opengoods.jobs.dedup --dry-run
|
||||
python -m opengoods.jobs.dedup --actor nightly
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import argparse
|
||||
import sys
|
||||
|
||||
import psycopg
|
||||
|
||||
from opengoods.etl.dedup import dedup_all
|
||||
from opengoods.etl.load import default_dsn
|
||||
|
||||
|
||||
def run(args: argparse.Namespace) -> int:
|
||||
with psycopg.connect(args.dsn, autocommit=False) as conn:
|
||||
summary = dedup_all(conn, actor=args.actor, dry_run=args.dry_run)
|
||||
if args.dry_run:
|
||||
conn.rollback()
|
||||
else:
|
||||
conn.commit()
|
||||
mode = "dry-run" if args.dry_run else "applied"
|
||||
print(f"{mode} groups={summary['groups']} merged={summary['merged']}")
|
||||
return 0
|
||||
|
||||
|
||||
def main(argv: list[str] | None = None) -> int:
|
||||
parser = argparse.ArgumentParser(description="Deduplicate OpenGoods products")
|
||||
parser.add_argument("--actor", default="ingestion", help="merge_log actor label")
|
||||
parser.add_argument("--dry-run", action="store_true", help="report only, do not write")
|
||||
parser.add_argument("--dsn", default=default_dsn(), help="PostgreSQL DSN")
|
||||
return run(parser.parse_args(argv))
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
sys.exit(main())
|
||||
@@ -0,0 +1,66 @@
|
||||
"""Lightweight recurring ingestion scheduler.
|
||||
|
||||
Runs one ingestion cycle (incremental OFF update, then dedup) on a fixed
|
||||
interval. Dependency-free: a plain sleep loop rather than a cron/APScheduler
|
||||
dependency, so it is trivial to run in a container or under systemd/supervisor.
|
||||
|
||||
Usage:
|
||||
python -m opengoods.jobs.schedule --once # single cycle, then exit
|
||||
python -m opengoods.jobs.schedule --interval 3600 # every hour
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import argparse
|
||||
import sys
|
||||
import time
|
||||
from datetime import UTC, datetime
|
||||
|
||||
from opengoods.etl.load import default_dsn
|
||||
from opengoods.jobs import dedup as dedup_job
|
||||
from opengoods.jobs import update_off as update_job
|
||||
|
||||
|
||||
def _cycle(args: argparse.Namespace) -> None:
|
||||
ts = datetime.now(UTC).isoformat(timespec="seconds")
|
||||
print(f"[{ts}] cycle start")
|
||||
update_job.run(
|
||||
argparse.Namespace(
|
||||
since=None,
|
||||
page_size=args.page_size,
|
||||
max_pages=args.max_pages,
|
||||
min_interval=args.min_interval,
|
||||
dsn=args.dsn,
|
||||
)
|
||||
)
|
||||
if not args.skip_dedup:
|
||||
dedup_job.run(argparse.Namespace(actor="scheduler", dry_run=False, dsn=args.dsn))
|
||||
|
||||
|
||||
def run(args: argparse.Namespace) -> int:
|
||||
_cycle(args)
|
||||
if args.once:
|
||||
return 0
|
||||
while True:
|
||||
time.sleep(args.interval)
|
||||
try:
|
||||
_cycle(args)
|
||||
except Exception as exc: # noqa: BLE001 - keep the loop alive across failures
|
||||
print(f"cycle error: {exc}", file=sys.stderr)
|
||||
return 0
|
||||
|
||||
|
||||
def main(argv: list[str] | None = None) -> int:
|
||||
parser = argparse.ArgumentParser(description="Recurring OpenGoods ingestion")
|
||||
parser.add_argument("--interval", type=int, default=3600, help="seconds between cycles")
|
||||
parser.add_argument("--once", action="store_true", help="run a single cycle and exit")
|
||||
parser.add_argument("--skip-dedup", action="store_true", help="run update only")
|
||||
parser.add_argument("--page-size", type=int, default=100, help="search page size")
|
||||
parser.add_argument("--max-pages", type=int, default=10, help="max pages to scan")
|
||||
parser.add_argument("--min-interval", type=float, default=4.0, help="API throttle seconds")
|
||||
parser.add_argument("--dsn", default=default_dsn(), help="PostgreSQL DSN")
|
||||
return run(parser.parse_args(argv))
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
sys.exit(main())
|
||||
@@ -0,0 +1,65 @@
|
||||
"""Incremental Open Food Facts update.
|
||||
|
||||
Fetches products modified since the persisted watermark, loads them, then
|
||||
advances the watermark to the newest ``last_modified_t`` processed so the next
|
||||
run only sees what changed.
|
||||
|
||||
Usage:
|
||||
python -m opengoods.jobs.update_off --max-pages 5
|
||||
python -m opengoods.jobs.update_off --since 1700000000 # override watermark
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import argparse
|
||||
import sys
|
||||
|
||||
import psycopg
|
||||
|
||||
from opengoods.adapters.openfoodfacts import SOURCE_NAME, OpenFoodFactsAdapter
|
||||
from opengoods.etl.load import default_dsn, ensure_source, load_record
|
||||
from opengoods.etl.state import get_watermark, set_watermark
|
||||
from opengoods.etl.transform import transform
|
||||
|
||||
|
||||
def run(args: argparse.Namespace) -> int:
|
||||
adapter = OpenFoodFactsAdapter(min_interval=args.min_interval)
|
||||
loaded = skipped = 0
|
||||
high_watermark = 0
|
||||
with psycopg.connect(args.dsn, autocommit=False) as conn:
|
||||
source_id = ensure_source(conn)
|
||||
since = args.since if args.since is not None else get_watermark(conn, SOURCE_NAME)
|
||||
high_watermark = since
|
||||
for raw in adapter.fetch_modified_since(
|
||||
since, page_size=args.page_size, max_pages=args.max_pages
|
||||
):
|
||||
high_watermark = max(high_watermark, int(raw.get("last_modified_t") or 0))
|
||||
rec = transform(raw)
|
||||
if rec is None:
|
||||
skipped += 1
|
||||
continue
|
||||
load_record(conn, rec, source_id, raw)
|
||||
loaded += 1
|
||||
set_watermark(
|
||||
conn,
|
||||
SOURCE_NAME,
|
||||
high_watermark,
|
||||
stats={"loaded": loaded, "skipped": skipped, "since": since},
|
||||
)
|
||||
conn.commit()
|
||||
print(f"since={since} loaded={loaded} skipped={skipped} watermark={high_watermark}")
|
||||
return 0
|
||||
|
||||
|
||||
def main(argv: list[str] | None = None) -> int:
|
||||
parser = argparse.ArgumentParser(description="Incremental OFF update")
|
||||
parser.add_argument("--since", type=int, default=None, help="override watermark (unix ts)")
|
||||
parser.add_argument("--page-size", type=int, default=100, help="search page size")
|
||||
parser.add_argument("--max-pages", type=int, default=10, help="max pages to scan")
|
||||
parser.add_argument("--min-interval", type=float, default=4.0, help="API throttle seconds")
|
||||
parser.add_argument("--dsn", default=default_dsn(), help="PostgreSQL DSN")
|
||||
return run(parser.parse_args(argv))
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
sys.exit(main())
|
||||
Reference in New Issue
Block a user