391 lines
13 KiB
Python
391 lines
13 KiB
Python
"""Load transformed product records into the OpenGoods PostgreSQL database.
|
|
|
|
Only the ingestion side writes to the database. Every load records OFF as the
|
|
source with field-level provenance in `product_source`.
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import json
|
|
import logging
|
|
import os
|
|
from typing import Any
|
|
|
|
import psycopg
|
|
from psycopg.types.json import Jsonb
|
|
|
|
from opengoods.adapters.bypos import SOURCE_HOMEPAGE as SOURCE_HOMEPAGE_BYPOS
|
|
from opengoods.adapters.bypos import SOURCE_LICENSE as SOURCE_LICENSE_BYPOS
|
|
from opengoods.adapters.bypos import SOURCE_NAME as SOURCE_NAME_BYPOS
|
|
from opengoods.adapters.bypos import SOURCE_TRUST as SOURCE_TRUST_BYPOS
|
|
from opengoods.adapters.openfoodfacts import OFF_LICENSE, SOURCE_NAME
|
|
from opengoods.etl.quality import update_quality
|
|
|
|
OFF_HOMEPAGE = "https://world.openfoodfacts.org"
|
|
|
|
logger = logging.getLogger(__name__)
|
|
|
|
|
|
def default_dsn() -> str:
|
|
return os.environ.get(
|
|
"OPENGOODS_DATABASE_URL",
|
|
"postgres://opengoods:opengoods@localhost:5432/opengoods?sslmode=disable",
|
|
)
|
|
|
|
|
|
def _normalize_brand(name: str) -> str:
|
|
return " ".join(name.lower().split())
|
|
|
|
|
|
def ensure_source_named(
|
|
conn: psycopg.Connection,
|
|
name: str,
|
|
homepage: str,
|
|
license: str,
|
|
trust_weight: float,
|
|
) -> str:
|
|
"""Upsert a source row by name and return its id."""
|
|
row = conn.execute(
|
|
"""
|
|
INSERT INTO source (name, homepage, license, trust_weight)
|
|
VALUES (%s, %s, %s, %s)
|
|
ON CONFLICT (name) DO UPDATE SET homepage = EXCLUDED.homepage
|
|
RETURNING id
|
|
""",
|
|
(name, homepage, license, trust_weight),
|
|
).fetchone()
|
|
return row[0]
|
|
|
|
|
|
def ensure_source(conn: psycopg.Connection) -> str:
|
|
"""Upsert the Open Food Facts source row and return its id."""
|
|
return ensure_source_named(conn, SOURCE_NAME, OFF_HOMEPAGE, OFF_LICENSE, 0.7)
|
|
|
|
|
|
def _ensure_brand(conn: psycopg.Connection, name: str | None) -> str | None:
|
|
if not name:
|
|
return None
|
|
row = conn.execute(
|
|
"""
|
|
INSERT INTO brand (name, normalized_name)
|
|
VALUES (%s, %s)
|
|
ON CONFLICT (normalized_name) DO UPDATE SET name = brand.name
|
|
RETURNING id
|
|
""",
|
|
(name, _normalize_brand(name)),
|
|
).fetchone()
|
|
return row[0]
|
|
|
|
|
|
def _ensure_manufacturer(
|
|
conn: psycopg.Connection, name: str | None, country: str | None = None
|
|
) -> str | None:
|
|
if not name:
|
|
return None
|
|
row = conn.execute(
|
|
"""
|
|
INSERT INTO manufacturer (name, normalized_name, country)
|
|
VALUES (%s, %s, %s)
|
|
ON CONFLICT (normalized_name) DO UPDATE SET name = manufacturer.name
|
|
RETURNING id
|
|
""",
|
|
(name, _normalize_brand(name), country),
|
|
).fetchone()
|
|
return row[0]
|
|
|
|
|
|
def _category_id(conn: psycopg.Connection, path: str | None) -> tuple[str | None, str | None]:
|
|
if not path:
|
|
return None, None
|
|
row = conn.execute(
|
|
"SELECT id, gpc_brick_code FROM category WHERE path = %s::ltree", (path,)
|
|
).fetchone()
|
|
return (row[0], row[1]) if row else (None, None)
|
|
|
|
|
|
def load_record(conn: psycopg.Connection, rec: dict[str, Any], source_id: str, raw: dict) -> str:
|
|
"""Upsert one transformed record; return the product id."""
|
|
brand_id = _ensure_brand(conn, rec.get("brand"))
|
|
category_id, gpc_brick = _category_id(conn, rec.get("category_path"))
|
|
|
|
fields = ["name", "brand", "net_content", "category", "country_of_origin"]
|
|
|
|
if rec.get("gtin"):
|
|
prod = conn.execute(
|
|
"""
|
|
INSERT INTO product (gtin, name, brand_id, category_id, gpc_brick_code,
|
|
net_content_value, net_content_unit, net_content_canonical,
|
|
country_of_origin, attributes)
|
|
VALUES (%s,%s,%s,%s,%s,%s,%s,%s,%s,%s)
|
|
ON CONFLICT (gtin) WHERE gtin IS NOT NULL DO UPDATE SET
|
|
name = EXCLUDED.name,
|
|
brand_id = COALESCE(EXCLUDED.brand_id, product.brand_id),
|
|
category_id = COALESCE(EXCLUDED.category_id, product.category_id),
|
|
gpc_brick_code = COALESCE(EXCLUDED.gpc_brick_code, product.gpc_brick_code),
|
|
net_content_value = EXCLUDED.net_content_value,
|
|
net_content_unit = EXCLUDED.net_content_unit,
|
|
net_content_canonical = EXCLUDED.net_content_canonical,
|
|
country_of_origin = EXCLUDED.country_of_origin
|
|
RETURNING id
|
|
""",
|
|
(
|
|
rec["gtin"],
|
|
rec["name"],
|
|
brand_id,
|
|
category_id,
|
|
gpc_brick,
|
|
rec.get("net_content_value"),
|
|
rec.get("net_content_unit"),
|
|
rec.get("net_content_canonical"),
|
|
rec.get("country_of_origin"),
|
|
Jsonb({}),
|
|
),
|
|
).fetchone()
|
|
else:
|
|
prod = conn.execute(
|
|
"""
|
|
INSERT INTO product (name, brand_id, category_id, gpc_brick_code,
|
|
net_content_value, net_content_unit, net_content_canonical,
|
|
country_of_origin, attributes)
|
|
VALUES (%s,%s,%s,%s,%s,%s,%s,%s,%s)
|
|
RETURNING id
|
|
""",
|
|
(
|
|
rec["name"],
|
|
brand_id,
|
|
category_id,
|
|
gpc_brick,
|
|
rec.get("net_content_value"),
|
|
rec.get("net_content_unit"),
|
|
rec.get("net_content_canonical"),
|
|
rec.get("country_of_origin"),
|
|
Jsonb({}),
|
|
),
|
|
).fetchone()
|
|
product_id = prod[0]
|
|
|
|
food = rec.get("food") or {}
|
|
conn.execute(
|
|
"""
|
|
INSERT INTO food_detail (product_id, ingredients_text, allergens, additives,
|
|
nutriments, nutrition_basis, serving_size, nutri_score)
|
|
VALUES (%s,%s,%s,%s,%s,%s,%s,%s)
|
|
ON CONFLICT (product_id) DO UPDATE SET
|
|
ingredients_text = EXCLUDED.ingredients_text,
|
|
allergens = EXCLUDED.allergens,
|
|
additives = EXCLUDED.additives,
|
|
nutriments = EXCLUDED.nutriments,
|
|
nutrition_basis = EXCLUDED.nutrition_basis,
|
|
serving_size = EXCLUDED.serving_size,
|
|
nutri_score = EXCLUDED.nutri_score
|
|
""",
|
|
(
|
|
product_id,
|
|
food.get("ingredients_text"),
|
|
food.get("allergens") or [],
|
|
food.get("additives") or [],
|
|
Jsonb(food.get("nutriments") or {}),
|
|
food.get("nutrition_basis"),
|
|
food.get("serving_size"),
|
|
food.get("nutri_score"),
|
|
),
|
|
)
|
|
|
|
if rec.get("image_url"):
|
|
conn.execute(
|
|
"""
|
|
INSERT INTO product_image (product_id, url, kind, license, source_id)
|
|
VALUES (%s,%s,'front',%s,%s)
|
|
""",
|
|
(product_id, rec["image_url"], "CC-BY-SA", source_id),
|
|
)
|
|
fields.append("image")
|
|
|
|
conn.execute(
|
|
"""
|
|
INSERT INTO product_source (product_id, source_id, url, fields, fetched_at, raw)
|
|
VALUES (%s,%s,%s,%s, now(), %s)
|
|
""",
|
|
(
|
|
product_id,
|
|
source_id,
|
|
f"{OFF_HOMEPAGE}/product/{rec.get('gtin') or ''}",
|
|
fields,
|
|
Jsonb(_jsonable(raw)),
|
|
),
|
|
)
|
|
|
|
# Recompute the data-quality score now that all facts + provenance exist.
|
|
update_quality(conn, product_id)
|
|
return product_id
|
|
|
|
|
|
def load_record_safe(
|
|
conn: psycopg.Connection, rec: dict[str, Any], source_id: str, raw: dict
|
|
) -> bool:
|
|
"""Load one record inside a savepoint.
|
|
|
|
On success the record's writes stay in the surrounding transaction. On any
|
|
error, only this record's writes are rolled back (to the savepoint) and the
|
|
batch continues, so a single malformed source record cannot abort a large
|
|
import. Returns True if loaded, False if skipped due to an error.
|
|
"""
|
|
try:
|
|
with conn.transaction():
|
|
load_record(conn, rec, source_id, raw)
|
|
return True
|
|
except Exception as exc: # noqa: BLE001 - per-record isolation is intentional
|
|
logger.warning("skipping record gtin=%s: %s", rec.get("gtin"), exc)
|
|
return False
|
|
|
|
|
|
def ensure_bypos_source(conn: psycopg.Connection) -> str:
|
|
"""Upsert the bypos central-library source row and return its id."""
|
|
return ensure_source_named(
|
|
conn,
|
|
SOURCE_NAME_BYPOS,
|
|
SOURCE_HOMEPAGE_BYPOS,
|
|
SOURCE_LICENSE_BYPOS,
|
|
SOURCE_TRUST_BYPOS,
|
|
)
|
|
|
|
|
|
def load_bypos_record(
|
|
conn: psycopg.Connection, rec: dict[str, Any], source_id: str, raw: dict
|
|
) -> str:
|
|
"""Upsert one transformed bypos record; return the product id.
|
|
|
|
Unlike OFF records these have no ingredients/nutrition, so no ``food_detail``
|
|
row is written. The suggested retail price (if any) is stored as a CNY MSRP
|
|
snapshot, and provenance/MSRP rows are keyed by source so a re-import
|
|
refreshes rather than duplicates them.
|
|
"""
|
|
manufacturer_id = _ensure_manufacturer(
|
|
conn, rec.get("manufacturer"), rec.get("country_of_origin")
|
|
)
|
|
attrs = rec.get("attributes") or {}
|
|
fields = ["name", "net_content", "country_of_origin"]
|
|
if manufacturer_id:
|
|
fields.append("manufacturer")
|
|
if attrs:
|
|
fields.append("attributes")
|
|
|
|
if rec.get("gtin"):
|
|
fields.append("gtin")
|
|
prod = conn.execute(
|
|
"""
|
|
INSERT INTO product (gtin, name, manufacturer_id,
|
|
net_content_value, net_content_unit, net_content_canonical,
|
|
country_of_origin, attributes)
|
|
VALUES (%s,%s,%s,%s,%s,%s,%s,%s)
|
|
ON CONFLICT (gtin) WHERE gtin IS NOT NULL DO UPDATE SET
|
|
name = EXCLUDED.name,
|
|
manufacturer_id = COALESCE(EXCLUDED.manufacturer_id, product.manufacturer_id),
|
|
net_content_value = COALESCE(EXCLUDED.net_content_value, product.net_content_value),
|
|
net_content_unit = COALESCE(EXCLUDED.net_content_unit, product.net_content_unit),
|
|
net_content_canonical = COALESCE(
|
|
EXCLUDED.net_content_canonical, product.net_content_canonical),
|
|
country_of_origin = COALESCE(EXCLUDED.country_of_origin, product.country_of_origin),
|
|
attributes = product.attributes || EXCLUDED.attributes
|
|
RETURNING id
|
|
""",
|
|
(
|
|
rec["gtin"],
|
|
rec["name"],
|
|
manufacturer_id,
|
|
rec.get("net_content_value"),
|
|
rec.get("net_content_unit"),
|
|
rec.get("net_content_canonical"),
|
|
rec.get("country_of_origin"),
|
|
Jsonb(attrs),
|
|
),
|
|
).fetchone()
|
|
else:
|
|
prod = conn.execute(
|
|
"""
|
|
INSERT INTO product (name, manufacturer_id,
|
|
net_content_value, net_content_unit, net_content_canonical,
|
|
country_of_origin, attributes)
|
|
VALUES (%s,%s,%s,%s,%s,%s,%s)
|
|
RETURNING id
|
|
""",
|
|
(
|
|
rec["name"],
|
|
manufacturer_id,
|
|
rec.get("net_content_value"),
|
|
rec.get("net_content_unit"),
|
|
rec.get("net_content_canonical"),
|
|
rec.get("country_of_origin"),
|
|
Jsonb(attrs),
|
|
),
|
|
).fetchone()
|
|
product_id = prod[0]
|
|
|
|
# Refresh this source's MSRP snapshot (suggested retail price, CNY).
|
|
conn.execute(
|
|
"DELETE FROM product_msrp WHERE product_id = %s AND source_id = %s",
|
|
(product_id, source_id),
|
|
)
|
|
if rec.get("msrp") is not None:
|
|
conn.execute(
|
|
"""
|
|
INSERT INTO product_msrp (product_id, amount, currency, region, source_id, source_url)
|
|
VALUES (%s,%s,'CNY','CN',%s,%s)
|
|
""",
|
|
(product_id, rec["msrp"], source_id, SOURCE_HOMEPAGE_BYPOS),
|
|
)
|
|
fields.append("msrp")
|
|
|
|
# Refresh this source's provenance row (one per source for idempotency).
|
|
conn.execute(
|
|
"DELETE FROM product_source WHERE product_id = %s AND source_id = %s",
|
|
(product_id, source_id),
|
|
)
|
|
conn.execute(
|
|
"""
|
|
INSERT INTO product_source (product_id, source_id, url, fields, fetched_at, raw)
|
|
VALUES (%s,%s,%s,%s, COALESCE(%s::timestamptz, now()), %s)
|
|
""",
|
|
(
|
|
product_id,
|
|
source_id,
|
|
SOURCE_HOMEPAGE_BYPOS,
|
|
fields,
|
|
rec.get("fetched_at"),
|
|
Jsonb(_jsonable(raw)),
|
|
),
|
|
)
|
|
|
|
update_quality(conn, product_id)
|
|
return product_id
|
|
|
|
|
|
def load_bypos_record_safe(
|
|
conn: psycopg.Connection, rec: dict[str, Any], source_id: str, raw: dict
|
|
) -> bool:
|
|
"""Load one bypos record inside a savepoint (see :func:`load_record_safe`)."""
|
|
try:
|
|
with conn.transaction():
|
|
load_bypos_record(conn, rec, source_id, raw)
|
|
return True
|
|
except Exception as exc: # noqa: BLE001 - per-record isolation is intentional
|
|
logger.warning("skipping bypos record gtin=%s: %s", rec.get("gtin"), exc)
|
|
return False
|
|
|
|
|
|
def _jsonable(raw: dict) -> dict:
|
|
"""Drop values that are not JSON-serializable from a raw record."""
|
|
try:
|
|
json.dumps(raw)
|
|
return raw
|
|
except (TypeError, ValueError):
|
|
return {k: v for k, v in raw.items() if _is_jsonable(v)}
|
|
|
|
|
|
def _is_jsonable(v: object) -> bool:
|
|
try:
|
|
json.dumps(v)
|
|
return True
|
|
except (TypeError, ValueError):
|
|
return False
|