feat(ingest): country-focused seeding (collect domestic CN products)
Add a market-focused seeding path so the catalogue can be built from domestic products rather than the English-heavy global default: - adapter.fetch_by_country(country): OFF search filtered by countries_tags_en, sorted by unique_scans_n (most-scanned first), de-duplicated across pages since OFF popularity ordering is unstable. - is_cn_gs1(code): True for GS1-China company prefixes (690-699), i.e. genuinely domestic items vs. imports merely sold in China. - seed_off --country <slug> [--domestic-only] [--page-size/--max-pages]: e.g. 'seed_off --country china --domestic-only' loads only 69x barcodes. Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>
This commit is contained in:
@@ -157,6 +157,60 @@ class OpenFoodFactsAdapter:
|
|||||||
if reached_old or len(products) < page_size:
|
if reached_old or len(products) < page_size:
|
||||||
return
|
return
|
||||||
|
|
||||||
|
def fetch_by_country(
|
||||||
|
self,
|
||||||
|
country: str,
|
||||||
|
*,
|
||||||
|
page_size: int = 100,
|
||||||
|
max_pages: int = 10,
|
||||||
|
sort_by: str = "unique_scans_n",
|
||||||
|
) -> Iterator[dict]:
|
||||||
|
"""Yield products sold in ``country`` (an OFF ``countries_tags_en`` slug).
|
||||||
|
|
||||||
|
Used to seed a market-specific catalogue (e.g. ``china``). Results are
|
||||||
|
sorted by ``sort_by`` (default ``unique_scans_n`` so the most-scanned,
|
||||||
|
best-known products come first) and de-duplicated across pages, since
|
||||||
|
OFF's popularity ordering is not stable between page requests.
|
||||||
|
"""
|
||||||
|
seen: set[str] = set()
|
||||||
|
for page in range(1, max_pages + 1):
|
||||||
|
resp = self._get(
|
||||||
|
_SEARCH_URL,
|
||||||
|
params={
|
||||||
|
"fields": _SEARCH_FIELDS,
|
||||||
|
"countries_tags_en": country,
|
||||||
|
"sort_by": sort_by,
|
||||||
|
"page": page,
|
||||||
|
"page_size": page_size,
|
||||||
|
},
|
||||||
|
)
|
||||||
|
products = resp.json().get("products") or []
|
||||||
|
if not products:
|
||||||
|
return
|
||||||
|
new_on_page = 0
|
||||||
|
for prod in products:
|
||||||
|
code = str(prod.get("code") or "")
|
||||||
|
if code and code in seen:
|
||||||
|
continue
|
||||||
|
if code:
|
||||||
|
seen.add(code)
|
||||||
|
new_on_page += 1
|
||||||
|
yield prod
|
||||||
|
if len(products) < page_size or new_on_page == 0:
|
||||||
|
return
|
||||||
|
|
||||||
|
|
||||||
|
def is_cn_gs1(code: str | None) -> bool:
|
||||||
|
"""Return True for a GS1 China company prefix (barcodes starting 690-699).
|
||||||
|
|
||||||
|
These identify products registered with GS1 China, i.e. genuinely domestic
|
||||||
|
items, as opposed to imported goods merely tagged as sold in China.
|
||||||
|
"""
|
||||||
|
if not code:
|
||||||
|
return False
|
||||||
|
code = code.strip()
|
||||||
|
return len(code) >= 3 and code[:2] == "69" and code[2].isdigit()
|
||||||
|
|
||||||
|
|
||||||
def read_dump(path: str | Path) -> Iterator[dict]:
|
def read_dump(path: str | Path) -> Iterator[dict]:
|
||||||
"""Yield raw product records from an OFF JSONL dump file.
|
"""Yield raw product records from an OFF JSONL dump file.
|
||||||
|
|||||||
@@ -7,6 +7,11 @@ Usage:
|
|||||||
# from a downloaded OFF JSONL dump (optionally .gz), limited to N records
|
# from a downloaded OFF JSONL dump (optionally .gz), limited to N records
|
||||||
python -m opengoods.jobs.seed_off --dump products.jsonl.gz --limit 1000
|
python -m opengoods.jobs.seed_off --dump products.jsonl.gz --limit 1000
|
||||||
|
|
||||||
|
# market-focused: the most-scanned products sold in China, restricted to
|
||||||
|
# genuine GS1-China (69x) barcodes
|
||||||
|
python -m opengoods.jobs.seed_off --country china --domestic-only \
|
||||||
|
--max-pages 20 --limit 1000
|
||||||
|
|
||||||
The OFF read API is rate-limited client-side; for large imports use a dump.
|
The OFF read API is rate-limited client-side; for large imports use a dump.
|
||||||
"""
|
"""
|
||||||
|
|
||||||
@@ -18,20 +23,33 @@ from collections.abc import Iterator
|
|||||||
|
|
||||||
import psycopg
|
import psycopg
|
||||||
|
|
||||||
from opengoods.adapters.openfoodfacts import OpenFoodFactsAdapter, read_dump
|
from opengoods.adapters.openfoodfacts import (
|
||||||
|
OpenFoodFactsAdapter,
|
||||||
|
is_cn_gs1,
|
||||||
|
read_dump,
|
||||||
|
)
|
||||||
from opengoods.etl.load import default_dsn, ensure_source, load_record_safe
|
from opengoods.etl.load import default_dsn, ensure_source, load_record_safe
|
||||||
from opengoods.etl.transform import transform
|
from opengoods.etl.transform import transform
|
||||||
|
|
||||||
|
|
||||||
def _raw_records(args: argparse.Namespace) -> Iterator[dict]:
|
def _raw_records(args: argparse.Namespace) -> Iterator[dict]:
|
||||||
if args.dump:
|
if args.dump:
|
||||||
records = read_dump(args.dump)
|
records: Iterator[dict] = read_dump(args.dump)
|
||||||
|
elif args.country:
|
||||||
|
adapter = OpenFoodFactsAdapter(min_interval=args.min_interval)
|
||||||
|
records = adapter.fetch_by_country(
|
||||||
|
args.country, page_size=args.page_size, max_pages=args.max_pages
|
||||||
|
)
|
||||||
else:
|
else:
|
||||||
adapter = OpenFoodFactsAdapter(min_interval=args.min_interval)
|
adapter = OpenFoodFactsAdapter(min_interval=args.min_interval)
|
||||||
records = adapter.fetch(args.barcodes)
|
records = adapter.fetch(args.barcodes)
|
||||||
for i, rec in enumerate(records):
|
yielded = 0
|
||||||
if args.limit and i >= args.limit:
|
for rec in records:
|
||||||
|
if args.domestic_only and not is_cn_gs1(str(rec.get("code") or "")):
|
||||||
|
continue
|
||||||
|
if args.limit and yielded >= args.limit:
|
||||||
break
|
break
|
||||||
|
yielded += 1
|
||||||
yield rec
|
yield rec
|
||||||
|
|
||||||
|
|
||||||
@@ -58,7 +76,18 @@ def main(argv: list[str] | None = None) -> int:
|
|||||||
src = parser.add_mutually_exclusive_group(required=True)
|
src = parser.add_mutually_exclusive_group(required=True)
|
||||||
src.add_argument("--barcodes", nargs="+", help="barcodes to fetch via the OFF API")
|
src.add_argument("--barcodes", nargs="+", help="barcodes to fetch via the OFF API")
|
||||||
src.add_argument("--dump", help="path to an OFF JSONL dump (.jsonl or .jsonl.gz)")
|
src.add_argument("--dump", help="path to an OFF JSONL dump (.jsonl or .jsonl.gz)")
|
||||||
|
src.add_argument(
|
||||||
|
"--country",
|
||||||
|
help="OFF countries_tags_en slug to seed from, e.g. 'china' (most-scanned first)",
|
||||||
|
)
|
||||||
|
parser.add_argument(
|
||||||
|
"--domestic-only",
|
||||||
|
action="store_true",
|
||||||
|
help="keep only genuine GS1-China (69x) barcodes; drop imported goods",
|
||||||
|
)
|
||||||
parser.add_argument("--limit", type=int, default=0, help="max records to load (0 = all)")
|
parser.add_argument("--limit", type=int, default=0, help="max records to load (0 = all)")
|
||||||
|
parser.add_argument("--page-size", type=int, default=100, help="search page size")
|
||||||
|
parser.add_argument("--max-pages", type=int, default=10, help="max search pages (country mode)")
|
||||||
parser.add_argument("--min-interval", type=float, default=4.0, help="API throttle seconds")
|
parser.add_argument("--min-interval", type=float, default=4.0, help="API throttle seconds")
|
||||||
parser.add_argument("--dsn", default=default_dsn(), help="PostgreSQL DSN")
|
parser.add_argument("--dsn", default=default_dsn(), help="PostgreSQL DSN")
|
||||||
return run(parser.parse_args(argv))
|
return run(parser.parse_args(argv))
|
||||||
|
|||||||
@@ -0,0 +1,57 @@
|
|||||||
|
import httpx
|
||||||
|
|
||||||
|
from opengoods.adapters.openfoodfacts import OpenFoodFactsAdapter, is_cn_gs1
|
||||||
|
|
||||||
|
|
||||||
|
def _adapter(handler, **kwargs):
|
||||||
|
client = httpx.Client(transport=httpx.MockTransport(handler))
|
||||||
|
return OpenFoodFactsAdapter(client=client, min_interval=0, backoff_base=0, **kwargs)
|
||||||
|
|
||||||
|
|
||||||
|
def test_is_cn_gs1():
|
||||||
|
assert is_cn_gs1("6901234567892") # GS1 China prefix
|
||||||
|
assert is_cn_gs1("690")
|
||||||
|
assert not is_cn_gs1("3017624010701") # France
|
||||||
|
assert not is_cn_gs1("5449000000996") # Belgium
|
||||||
|
assert not is_cn_gs1("")
|
||||||
|
assert not is_cn_gs1(None)
|
||||||
|
assert not is_cn_gs1("69") # too short to carry a prefix digit
|
||||||
|
|
||||||
|
|
||||||
|
def test_fetch_by_country_passes_filter_and_paginates():
|
||||||
|
seen_params = []
|
||||||
|
|
||||||
|
def handler(request: httpx.Request) -> httpx.Response:
|
||||||
|
seen_params.append(dict(request.url.params))
|
||||||
|
page = int(request.url.params.get("page", "1"))
|
||||||
|
if page == 1:
|
||||||
|
return httpx.Response(
|
||||||
|
200,
|
||||||
|
json={"products": [{"code": "6901"}, {"code": "6902"}]},
|
||||||
|
)
|
||||||
|
return httpx.Response(200, json={"products": []})
|
||||||
|
|
||||||
|
adapter = _adapter(handler)
|
||||||
|
out = list(adapter.fetch_by_country("china", page_size=2, max_pages=5))
|
||||||
|
assert [p["code"] for p in out] == ["6901", "6902"]
|
||||||
|
assert seen_params[0]["countries_tags_en"] == "china"
|
||||||
|
assert seen_params[0]["sort_by"] == "unique_scans_n"
|
||||||
|
|
||||||
|
|
||||||
|
def test_fetch_by_country_dedupes_across_pages():
|
||||||
|
def handler(request: httpx.Request) -> httpx.Response:
|
||||||
|
page = int(request.url.params.get("page", "1"))
|
||||||
|
if page == 1:
|
||||||
|
return httpx.Response(
|
||||||
|
200, json={"products": [{"code": "6901"}, {"code": "6902"}]}
|
||||||
|
)
|
||||||
|
if page == 2:
|
||||||
|
# OFF popularity ordering is unstable; a repeat appears on page 2
|
||||||
|
return httpx.Response(
|
||||||
|
200, json={"products": [{"code": "6902"}, {"code": "6903"}]}
|
||||||
|
)
|
||||||
|
return httpx.Response(200, json={"products": []})
|
||||||
|
|
||||||
|
adapter = _adapter(handler)
|
||||||
|
out = [p["code"] for p in adapter.fetch_by_country("china", page_size=2, max_pages=5)]
|
||||||
|
assert out == ["6901", "6902", "6903"]
|
||||||
Reference in New Issue
Block a user