feat(ingestion): harden OFF ingestion for bulk seeding
- Add retry/backoff (429 + 5xx, Retry-After aware) to the OFF adapter so transient API errors no longer abort a run. - Clamp bounded text fields (serving_size, net_content_unit, country_of_origin) to their column widths in transform; long OFF values previously raised StringDataRightTruncation and rolled back the batch. - Load each record inside a savepoint (load_record_safe) so one malformed source record is skipped instead of aborting the whole import; jobs now report an errored count. - Tests for retry behaviour, serving_size clamping, and per-record isolation. Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>
This commit is contained in:
@@ -19,7 +19,7 @@ from collections.abc import Iterator
|
||||
import psycopg
|
||||
|
||||
from opengoods.adapters.openfoodfacts import OpenFoodFactsAdapter, read_dump
|
||||
from opengoods.etl.load import default_dsn, ensure_source, load_record
|
||||
from opengoods.etl.load import default_dsn, ensure_source, load_record_safe
|
||||
from opengoods.etl.transform import transform
|
||||
|
||||
|
||||
@@ -36,7 +36,7 @@ def _raw_records(args: argparse.Namespace) -> Iterator[dict]:
|
||||
|
||||
|
||||
def run(args: argparse.Namespace) -> int:
|
||||
loaded = skipped = 0
|
||||
loaded = skipped = errored = 0
|
||||
with psycopg.connect(args.dsn, autocommit=False) as conn:
|
||||
source_id = ensure_source(conn)
|
||||
for raw in _raw_records(args):
|
||||
@@ -44,10 +44,12 @@ def run(args: argparse.Namespace) -> int:
|
||||
if rec is None:
|
||||
skipped += 1
|
||||
continue
|
||||
load_record(conn, rec, source_id, raw)
|
||||
loaded += 1
|
||||
if load_record_safe(conn, rec, source_id, raw):
|
||||
loaded += 1
|
||||
else:
|
||||
errored += 1
|
||||
conn.commit()
|
||||
print(f"loaded={loaded} skipped={skipped}")
|
||||
print(f"loaded={loaded} skipped={skipped} errored={errored}")
|
||||
return 0
|
||||
|
||||
|
||||
|
||||
@@ -17,14 +17,14 @@ 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.load import default_dsn, ensure_source, load_record_safe
|
||||
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
|
||||
loaded = skipped = errored = 0
|
||||
high_watermark = 0
|
||||
with psycopg.connect(args.dsn, autocommit=False) as conn:
|
||||
source_id = ensure_source(conn)
|
||||
@@ -38,16 +38,21 @@ def run(args: argparse.Namespace) -> int:
|
||||
if rec is None:
|
||||
skipped += 1
|
||||
continue
|
||||
load_record(conn, rec, source_id, raw)
|
||||
loaded += 1
|
||||
if load_record_safe(conn, rec, source_id, raw):
|
||||
loaded += 1
|
||||
else:
|
||||
errored += 1
|
||||
set_watermark(
|
||||
conn,
|
||||
SOURCE_NAME,
|
||||
high_watermark,
|
||||
stats={"loaded": loaded, "skipped": skipped, "since": since},
|
||||
stats={"loaded": loaded, "skipped": skipped, "errored": errored, "since": since},
|
||||
)
|
||||
conn.commit()
|
||||
print(f"since={since} loaded={loaded} skipped={skipped} watermark={high_watermark}")
|
||||
print(
|
||||
f"since={since} loaded={loaded} skipped={skipped} "
|
||||
f"errored={errored} watermark={high_watermark}"
|
||||
)
|
||||
return 0
|
||||
|
||||
|
||||
|
||||
Reference in New Issue
Block a user