Compare commits

..

2 Commits

Author SHA1 Message Date
lixu 766a573989 feat(M2): Open Food Facts 采集导入 ETL
- adapters/openfoodfacts.py: OFF API 适配器(限速+User-Agent) + JSONL/.gz dump 读取
- etl/transform.py: 纯函数转换(GTIN 校验/净含量解析+归一/营养 per_100g 双能量/过敏原添加剂清洗/关键词分类映射)
- etl/load.py: psycopg upsert(product/food_detail/product_image) + product_source 字段级溯源
- jobs/seed_off.py: CLI(--barcodes API / --dump 文件 / --limit)
- 测试: 13 个离线 transform 单测(fixture) + 可跳过的 DB 集成测试
- 依赖: 增加 psycopg[binary]
- 实跑: 从 OFF API 拉 5 个真实条码入本地 postgres 验证通过
- docs/etl-openfoodfacts.md: 流程/运行/字段映射

Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>
2026-06-08 07:01:04 +00:00
lixu b0b816b0ee feat(M1): 数据模型迁移 + GS1 GPC 食品分类 + 单位字典
- migrations/0001_init: 全部核心表(product/food_detail/product_msrp/product_source/brand/manufacturer/category/category_schema/unit/attribute_definition/product_image/merge_log) + 索引(gtin唯一/name trigram/JSONB GIN/category ltree/tsvector) + tsvector/updated_at 触发器
- 0002_seed_units: 单位字典(与 units.py 一致, 含中文别名) + 常用营养参数定义
- 0003_seed_categories: 食品品类骨架(GS1 GPC 映射 + 自建中文树, ltree) + 品类参数模板(营养基准 per_100g/ml)
- CI 增加 migrations job: 用 postgres service 跑 migrate up + down 验证可逆
- 本地实跑 up/down/re-up 全部通过

Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>
2026-06-08 06:32:16 +00:00
17 changed files with 1047 additions and 1 deletions
+29
View File
@@ -43,3 +43,32 @@ jobs:
run: ruff format --check . run: ruff format --check .
- name: Pytest - name: Pytest
run: pytest -q run: pytest -q
migrations:
name: Migrations (postgres)
runs-on: ubuntu-latest
services:
postgres:
image: postgres:16-alpine
env:
POSTGRES_USER: opengoods
POSTGRES_PASSWORD: opengoods
POSTGRES_DB: opengoods
ports:
- "5432:5432"
options: >-
--health-cmd "pg_isready -U opengoods"
--health-interval 5s --health-timeout 5s --health-retries 10
env:
DBURL: postgres://opengoods:opengoods@localhost:5432/opengoods?sslmode=disable
steps:
- uses: actions/checkout@v4
- uses: actions/setup-go@v5
with:
go-version: "1.23"
- name: Install golang-migrate
run: go install -tags 'postgres' github.com/golang-migrate/migrate/v4/cmd/migrate@v4.18.1
- name: Migrate up
run: migrate -path migrations -database "$DBURL" up
- name: Migrate down (reversibility)
run: migrate -path migrations -database "$DBURL" down -all
+38
View File
@@ -0,0 +1,38 @@
# ETL: Open Food Facts 导入 (M2)
把 Open Food Facts (OFF, ODbL 许可) 的食品数据采集、转换并入库。只有 Python 采集侧写库,每条记录都以 `openfoodfacts` 为来源记录**字段级溯源**。
## 流程
```
OFF API / dump(jsonl[.gz])
→ adapters/openfoodfacts.py # 读取(限速 + User-Agent) / 解析 dump
→ etl/transform.py # 字段映射 + 单位归一 + 营养 per_100g + 分类映射(关键词)
→ etl/load.py # psycopg upsert(product/food_detail/product_image) + product_source 溯源
```
## 运行
先确保本地依赖与迁移就绪:`docker compose up -d postgres` + `migrate ... up`
```bash
# 用 OFF API 拉指定条码(客户端限速, 默认 4s/次)
python -m opengoods.jobs.seed_off --barcodes 3017624010701 5449000000996
# 用下载好的 OFF dump 批量导入(可 .gz), 限制条数
python -m opengoods.jobs.seed_off --dump products.jsonl.gz --limit 1000
```
DSN 默认读 `OPENGOODS_DATABASE_URL`
## 字段映射要点
| OFF | OpenGoods | 处理 |
|-----|-----------|------|
| `code` | `product.gtin` | GTIN-8/12/13/14 校验位验证, 不合法则不作为 gtin |
| `product_name_zh/_/_en` | `product.name` | 优先中文 |
| `brands` | `brand` | 取第一个, normalized_name 去重 |
| `quantity` | `net_content_*` | 解析 "500 g"/"1,5 L" → 经 `units.py` 归一(原始+归一双存) |
| `nutriments.*_100g` | `food_detail.nutriments` | per_100g; 能量 kJ/kcal 双存, 缺一自动换算 |
| `allergens_tags`/`additives_tags` | `allergens`/`additives` | 去 `en:` 前缀 |
| `nutriscore_grade` | `nutri_score` | 大写单字母 |
| `categories*`/name | `category_id` | 关键词映射到自建品类树(起步版, 后续换 OFF 分类→GPC 映射表) |
| `image_front_url` | `product_image` | 标 CC-BY-SA 许可 |
> 全量 dump 约数 GBCI 与单测用 fixture 离线验证 transform,DB 集成测试在无库时自动跳过。
@@ -0,0 +1,86 @@
"""Open Food Facts (OFF) source adapter.
Fetches raw product records either from the OFF read API (one product per
barcode) or from a downloaded JSONL dump file. OFF data is licensed under the
Open Database License (ODbL); product images are CC-BY-SA. We record OFF as the
source for every field we ingest.
The adapter is read-only and rate-limited to stay well within OFF's API limits
(<= ~15 req/min/IP for product reads) and to be a good citizen.
"""
from __future__ import annotations
import json
import time
from collections.abc import Iterator
from pathlib import Path
import httpx
SOURCE_NAME = "openfoodfacts"
OFF_LICENSE = "ODbL"
USER_AGENT = "OpenGoods/0.1 (+https://github.com/baicai2026-baicai/goods) public-good product API"
# Conservative client-side spacing between API calls (seconds).
_DEFAULT_MIN_INTERVAL = 4.0
_API_URL = "https://world.openfoodfacts.org/api/v2/product/{barcode}.json"
class OpenFoodFactsAdapter:
"""Read product records from the OFF API."""
source_name = SOURCE_NAME
def __init__(
self,
client: httpx.Client | None = None,
min_interval: float = _DEFAULT_MIN_INTERVAL,
) -> None:
self._client = client or httpx.Client(headers={"User-Agent": USER_AGENT}, timeout=30.0)
self._min_interval = min_interval
self._last_call = 0.0
def _throttle(self) -> None:
elapsed = time.monotonic() - self._last_call
wait = self._min_interval - elapsed
if wait > 0:
time.sleep(wait)
self._last_call = time.monotonic()
def fetch_barcode(self, barcode: str) -> dict | None:
"""Fetch a single product by barcode; return the raw `product` dict."""
self._throttle()
resp = self._client.get(_API_URL.format(barcode=barcode))
resp.raise_for_status()
payload = resp.json()
if payload.get("status") != 1:
return None
return payload["product"]
def fetch(self, barcodes: list[str]) -> Iterator[dict]:
"""Yield raw product records for the given barcodes."""
for code in barcodes:
record = self.fetch_barcode(code)
if record is not None:
yield record
def read_dump(path: str | Path) -> Iterator[dict]:
"""Yield raw product records from an OFF JSONL dump file.
Each line is one product JSON object (the format of OFF's .jsonl export).
Supports plain or .gz files.
"""
p = Path(path)
if p.suffix == ".gz":
import gzip
opener = lambda: gzip.open(p, "rt", encoding="utf-8") # noqa: E731
else:
opener = lambda: open(p, encoding="utf-8") # noqa: E731
with opener() as fh:
for line in fh:
line = line.strip()
if line:
yield json.loads(line)
+198
View File
@@ -0,0 +1,198 @@
"""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 os
from typing import Any
import psycopg
from psycopg.types.json import Jsonb
from opengoods.adapters.openfoodfacts import OFF_LICENSE, SOURCE_NAME
OFF_HOMEPAGE = "https://world.openfoodfacts.org"
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(conn: psycopg.Connection) -> str:
"""Upsert the Open Food Facts source row 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
""",
(SOURCE_NAME, OFF_HOMEPAGE, OFF_LICENSE, 0.7),
).fetchone()
return row[0]
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 _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)),
),
)
return product_id
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
+159
View File
@@ -0,0 +1,159 @@
"""Transform raw Open Food Facts records into the OpenGoods internal shape.
Pure functions (no DB, no network) so they are easy to unit-test against
fixtures. The output dict mirrors the columns the loader writes.
"""
from __future__ import annotations
import re
from decimal import Decimal
from opengoods.units import UnitError, normalize
# OFF nutriment key -> our attribute key. Energy handled separately.
_NUTRIMENT_KEYS = {
"proteins_100g": "proteins",
"fat_100g": "fat",
"saturated-fat_100g": "saturated_fat",
"carbohydrates_100g": "carbohydrates",
"sugars_100g": "sugars",
"salt_100g": "salt",
}
# Very small keyword -> category path map (starter; replaced by a proper
# OFF taxonomy -> GPC mapping table later).
_CATEGORY_KEYWORDS: list[tuple[tuple[str, ...], str]] = [
(("water", "eau", "饮用水", "矿泉水"), "food.beverages.water"),
(("soda", "carbonated", "汽水", "碳酸"), "food.beverages.carbonated"),
(("juice", "jus", "果汁"), "food.beverages.juice"),
(("milk", "lait", "牛奶"), "food.dairy.milk"),
(("yogurt", "yoghurt", "yaourt", "酸奶"), "food.dairy.yogurt"),
(("cheese", "fromage", "奶酪", "干酪"), "food.dairy.cheese"),
(("bread", "pain", "面包"), "food.bakery.bread"),
(("biscuit", "cookie", "饼干"), "food.bakery.biscuits"),
(("chips", "crisps", "薯片", "膨化"), "food.snacks.chips"),
(("chocolate", "chocolat", "巧克力"), "food.snacks.chocolate"),
(("rice", "riz", "大米", "稻米"), "food.staple.rice"),
(("noodle", "pasta", "面条", "挂面"), "food.staple.noodles"),
(("oil", "huile", "食用油", "食油"), "food.staple.cooking_oil"),
(("soy sauce", "酱油"), "food.condiments.soy_sauce"),
(("salt", "sel", "食盐"), "food.condiments.salt"),
]
_QTY_RE = re.compile(r"(?P<value>\d+(?:[.,]\d+)?)\s*(?P<unit>[a-zA-Z\u4e00-\u9fff%]+)")
def is_valid_gtin(code: str) -> bool:
"""Validate a GTIN-8/12/13/14 using the standard check digit."""
if not code.isdigit() or len(code) not in (8, 12, 13, 14):
return False
digits = [int(c) for c in code]
check = digits[-1]
body = digits[:-1][::-1]
total = sum(d * (3 if i % 2 == 0 else 1) for i, d in enumerate(body))
return (10 - total % 10) % 10 == check
def parse_quantity(text: str) -> tuple[Decimal, str] | None:
"""Parse a free-text quantity like '500 g' or '1,5 L' -> (value, unit)."""
if not text:
return None
m = _QTY_RE.search(text)
if not m:
return None
value = Decimal(m.group("value").replace(",", "."))
return value, m.group("unit")
def map_category(raw: dict) -> str | None:
"""Best-effort map OFF categories/name to a self-built category path."""
haystack = " ".join(
str(raw.get(k, ""))
for k in ("categories", "categories_tags", "product_name", "product_name_en")
).lower()
for keywords, path in _CATEGORY_KEYWORDS:
if any(kw.lower() in haystack for kw in keywords):
return path
return None
def _clean_tags(tags: list[str] | None, prefix: str = "") -> list[str]:
out: list[str] = []
for t in tags or []:
v = t.split(":", 1)[-1] if ":" in t else t
v = v.strip().replace("-", " ")
if v:
out.append(v)
return out
def transform_nutriments(off_nutriments: dict) -> dict:
"""Build a nutriments dict on a per_100g basis with dual energy units."""
out: dict[str, object] = {}
for off_key, our_key in _NUTRIMENT_KEYS.items():
if off_key in off_nutriments and off_nutriments[off_key] is not None:
out[our_key] = float(off_nutriments[off_key])
kj = off_nutriments.get("energy-kj_100g")
kcal = off_nutriments.get("energy-kcal_100g")
if kj is None and kcal is not None:
kj = float(Decimal(str(kcal)) * Decimal("4.184"))
if kcal is None and kj is not None:
kcal = float(Decimal(str(kj)) / Decimal("4.184"))
if kj is not None:
out["energy_kj"] = round(float(kj), 3)
if kcal is not None:
out["energy_kcal"] = round(float(kcal), 3)
return out
def transform(raw: dict) -> dict | None:
"""Transform one raw OFF product record into an internal product dict.
Returns None if the record lacks a usable name.
"""
name = raw.get("product_name_zh") or raw.get("product_name") or raw.get("product_name_en")
if not name:
return None
code = str(raw.get("code", "")).strip()
gtin = code if code and is_valid_gtin(code) else None
brands = raw.get("brands") or ""
brand = brands.split(",")[0].strip() or None
net_value = net_unit = net_canonical = None
parsed = parse_quantity(raw.get("quantity", ""))
if parsed:
value, unit = parsed
try:
norm = normalize(value, unit)
net_value, net_unit, net_canonical = (
norm.value,
norm.unit,
norm.canonical_value,
)
except UnitError:
net_value, net_unit = value, unit
return {
"gtin": gtin,
"name": str(name).strip(),
"brand": brand,
"category_path": map_category(raw),
"net_content_value": net_value,
"net_content_unit": net_unit,
"net_content_canonical": net_canonical,
"country_of_origin": (raw.get("countries") or "").split(",")[0].strip() or None,
"food": {
"ingredients_text": raw.get("ingredients_text") or None,
"allergens": _clean_tags(raw.get("allergens_tags")),
"additives": _clean_tags(raw.get("additives_tags")),
"nutriments": transform_nutriments(raw.get("nutriments") or {}),
"nutrition_basis": "per_100g",
"serving_size": raw.get("serving_size") or None,
"nutri_score": (raw.get("nutriscore_grade") or "").upper()[:1] or None,
},
"image_url": raw.get("image_front_url") or raw.get("image_url") or None,
}
+66
View File
@@ -0,0 +1,66 @@
"""Seed the database with Open Food Facts data.
Usage:
# from a list of barcodes via the OFF API
python -m opengoods.jobs.seed_off --barcodes 3017624010701 5449000000996
# from a downloaded OFF JSONL dump (optionally .gz), limited to N records
python -m opengoods.jobs.seed_off --dump products.jsonl.gz --limit 1000
The OFF read API is rate-limited client-side; for large imports use a dump.
"""
from __future__ import annotations
import argparse
import sys
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.transform import transform
def _raw_records(args: argparse.Namespace) -> Iterator[dict]:
if args.dump:
records = read_dump(args.dump)
else:
adapter = OpenFoodFactsAdapter(min_interval=args.min_interval)
records = adapter.fetch(args.barcodes)
for i, rec in enumerate(records):
if args.limit and i >= args.limit:
break
yield rec
def run(args: argparse.Namespace) -> int:
loaded = skipped = 0
with psycopg.connect(args.dsn, autocommit=False) as conn:
source_id = ensure_source(conn)
for raw in _raw_records(args):
rec = transform(raw)
if rec is None:
skipped += 1
continue
load_record(conn, rec, source_id, raw)
loaded += 1
conn.commit()
print(f"loaded={loaded} skipped={skipped}")
return 0
def main(argv: list[str] | None = None) -> int:
parser = argparse.ArgumentParser(description="Seed OpenGoods from Open Food Facts")
src = parser.add_mutually_exclusive_group(required=True)
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)")
parser.add_argument("--limit", type=int, default=0, help="max records to load (0 = all)")
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())
+1
View File
@@ -5,6 +5,7 @@ description = "OpenGoods (天工·商品标签) ingestion & ETL: collect product
requires-python = ">=3.11" requires-python = ">=3.11"
dependencies = [ dependencies = [
"httpx>=0.27", "httpx>=0.27",
"psycopg[binary]>=3.2",
] ]
[project.optional-dependencies] [project.optional-dependencies]
+26
View File
@@ -0,0 +1,26 @@
{
"code": "3017624010701",
"product_name": "Nutella",
"product_name_en": "Nutella hazelnut spread",
"brands": "Ferrero, Nutella",
"quantity": "400 g",
"countries": "France, China",
"categories": "Spreads, Hazelnut spreads, Chocolate spreads",
"categories_tags": ["en:spreads", "en:chocolate-spreads"],
"ingredients_text": "Sugar, palm oil, hazelnuts, cocoa, skimmed milk powder",
"allergens_tags": ["en:milk", "en:nuts"],
"additives_tags": ["en:e322"],
"serving_size": "15 g",
"nutriscore_grade": "e",
"image_front_url": "https://images.openfoodfacts.org/images/products/301/762/401/0701/front_en.jpg",
"nutriments": {
"energy-kj_100g": 2252,
"energy-kcal_100g": 539,
"fat_100g": 30.9,
"saturated-fat_100g": 10.6,
"carbohydrates_100g": 57.5,
"sugars_100g": 56.3,
"proteins_100g": 6.3,
"salt_100g": 0.107
}
}
+61
View File
@@ -0,0 +1,61 @@
"""Integration test for the DB loader.
Skipped automatically when no database is reachable (e.g. local runs without
docker, or CI jobs without a postgres service). Requires migrations applied.
"""
from __future__ import annotations
import json
from pathlib import Path
import pytest
from opengoods.etl.load import default_dsn, ensure_source, load_record
from opengoods.etl.transform import transform
psycopg = pytest.importorskip("psycopg")
FIXTURE = json.loads((Path(__file__).parent / "fixtures" / "off_product.json").read_text())
@pytest.fixture()
def conn():
try:
c = psycopg.connect(default_dsn(), connect_timeout=3)
except psycopg.OperationalError as exc: # pragma: no cover - env dependent
pytest.skip(f"no database available: {exc}")
# ensure schema present
has_product = c.execute("SELECT to_regclass('public.product') IS NOT NULL").fetchone()[0]
if not has_product:
c.close()
pytest.skip("migrations not applied")
yield c
c.rollback()
c.close()
def test_load_record_roundtrip(conn):
source_id = ensure_source(conn)
rec = transform(FIXTURE)
product_id = load_record(conn, rec, source_id, FIXTURE)
row = conn.execute(
"SELECT name, gtin, net_content_unit FROM product WHERE id = %s", (product_id,)
).fetchone()
assert row[0] == "Nutella"
assert row[1] == "3017624010701"
assert row[2] == "g"
nutri = conn.execute(
"SELECT nutriments ->> 'energy_kcal' FROM food_detail WHERE product_id = %s",
(product_id,),
).fetchone()
assert nutri[0] == "539.0"
prov = conn.execute(
"SELECT count(*) FROM product_source WHERE product_id = %s", (product_id,)
).fetchone()
assert prov[0] >= 1
conn.rollback() # keep the test DB clean
+67
View File
@@ -0,0 +1,67 @@
import json
from decimal import Decimal
from pathlib import Path
import pytest
from opengoods.etl.transform import (
is_valid_gtin,
map_category,
parse_quantity,
transform,
transform_nutriments,
)
FIXTURE = json.loads((Path(__file__).parent / "fixtures" / "off_product.json").read_text())
def test_is_valid_gtin():
assert is_valid_gtin("3017624010701") # real EAN-13
assert is_valid_gtin("5449000000996") # Coca-Cola
assert not is_valid_gtin("3017624010700") # bad check digit
assert not is_valid_gtin("123")
assert not is_valid_gtin("notanumber")
def test_parse_quantity():
assert parse_quantity("400 g") == (Decimal("400"), "g")
assert parse_quantity("1,5 L") == (Decimal("1.5"), "L")
assert parse_quantity("") is None
assert parse_quantity("family size") is None
def test_map_category():
assert map_category({"product_name": "Spring Water"}) == "food.beverages.water"
assert map_category({"categories": "Dark chocolate"}) == "food.snacks.chocolate"
assert map_category({"product_name": "Mystery"}) is None
def test_transform_nutriments_dual_energy():
out = transform_nutriments(FIXTURE["nutriments"])
assert out["energy_kj"] == 2252.0
assert out["energy_kcal"] == 539.0
assert out["fat"] == 30.9
assert out["salt"] == 0.107
def test_transform_nutriments_fills_missing_energy():
out = transform_nutriments({"energy-kcal_100g": 100})
assert out["energy_kj"] == pytest.approx(418.4)
def test_transform_full_record():
rec = transform(FIXTURE)
assert rec is not None
assert rec["gtin"] == "3017624010701"
assert rec["name"] == "Nutella"
assert rec["brand"] == "Ferrero"
assert rec["net_content_unit"] == "g"
assert rec["net_content_canonical"] == Decimal("400")
assert rec["country_of_origin"] == "France"
assert rec["food"]["nutri_score"] == "E"
assert "milk" in rec["food"]["allergens"]
assert rec["image_url"].endswith(".jpg")
def test_transform_drops_unnamed():
assert transform({"code": "0000000000000"}) is None
+21
View File
@@ -0,0 +1,21 @@
DROP TRIGGER IF EXISTS trg_product_sync ON product;
DROP FUNCTION IF EXISTS product_sync_tsv();
DROP TABLE IF EXISTS merge_log;
DROP TABLE IF EXISTS product_source;
DROP TABLE IF EXISTS product_image;
DROP TABLE IF EXISTS product_msrp;
DROP TABLE IF EXISTS food_detail;
DROP TABLE IF EXISTS product;
DROP TABLE IF EXISTS attribute_definition;
DROP TABLE IF EXISTS unit;
DROP TABLE IF EXISTS category_schema;
DROP TABLE IF EXISTS category;
DROP TABLE IF EXISTS manufacturer;
DROP TABLE IF EXISTS brand;
DROP TABLE IF EXISTS source;
DROP EXTENSION IF EXISTS ltree;
DROP EXTENSION IF EXISTS pg_trgm;
-- keep pgcrypto (commonly shared); drop only if you are sure:
-- DROP EXTENSION IF EXISTS pgcrypto;
+182
View File
@@ -0,0 +1,182 @@
-- OpenGoods (天工·商品标签) initial schema.
-- Public-good product information store: facts only, no commerce.
CREATE EXTENSION IF NOT EXISTS pgcrypto; -- gen_random_uuid()
CREATE EXTENSION IF NOT EXISTS pg_trgm; -- fuzzy name search
CREATE EXTENSION IF NOT EXISTS ltree; -- category subtree queries
-- Data sources (Open Food Facts / USDA / GS1 ...) with trust + license.
CREATE TABLE source (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
name TEXT NOT NULL,
homepage TEXT,
license TEXT,
trust_weight NUMERIC(3,2) NOT NULL DEFAULT 0.5,
notes TEXT,
UNIQUE (name)
);
CREATE TABLE brand (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
name TEXT NOT NULL,
normalized_name TEXT NOT NULL,
aliases TEXT[] NOT NULL DEFAULT '{}',
UNIQUE (normalized_name)
);
CREATE TABLE manufacturer (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
name TEXT NOT NULL,
normalized_name TEXT NOT NULL,
country VARCHAR(64),
UNIQUE (normalized_name)
);
-- Self-built category tree, each node optionally mapped to a GS1 GPC brick.
CREATE TABLE category (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
name_zh TEXT NOT NULL,
name_en TEXT,
parent_id UUID REFERENCES category(id),
path LTREE NOT NULL,
gpc_brick_code VARCHAR(10),
level INT NOT NULL DEFAULT 0,
UNIQUE (path)
);
-- Parameter template / constraints per category.
CREATE TABLE category_schema (
category_id UUID PRIMARY KEY REFERENCES category(id) ON DELETE CASCADE,
required_attributes TEXT[] NOT NULL DEFAULT '{}',
recommended_attributes TEXT[] NOT NULL DEFAULT '{}',
nutriment_basis VARCHAR(16)
);
-- Unit dictionary: each unit maps to a canonical unit within its dimension.
CREATE TABLE unit (
code VARCHAR(16) PRIMARY KEY,
dimension VARCHAR(16) NOT NULL,
canonical VARCHAR(16) NOT NULL,
to_canonical_factor NUMERIC,
aliases TEXT[] NOT NULL DEFAULT '{}',
display TEXT
);
-- Parameter dictionary: standard attribute keys with default unit.
CREATE TABLE attribute_definition (
key VARCHAR(64) PRIMARY KEY,
label_zh TEXT,
label_en TEXT,
dimension VARCHAR(16),
default_unit VARCHAR(16) REFERENCES unit(code),
aliases TEXT[] NOT NULL DEFAULT '{}'
);
CREATE TABLE product (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
gtin VARCHAR(14),
name TEXT NOT NULL,
brand_id UUID REFERENCES brand(id),
manufacturer_id UUID REFERENCES manufacturer(id),
category_id UUID REFERENCES category(id),
gpc_brick_code VARCHAR(10),
net_content_value NUMERIC,
net_content_unit VARCHAR(16),
net_content_canonical NUMERIC,
country_of_origin VARCHAR(64),
shelf_life_days INT,
storage TEXT,
attributes JSONB NOT NULL DEFAULT '{}',
quality_score NUMERIC(4,3) NOT NULL DEFAULT 0,
status VARCHAR(16) NOT NULL DEFAULT 'active',
canonical_id UUID REFERENCES product(id),
search_tsv TSVECTOR,
created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
updated_at TIMESTAMPTZ NOT NULL DEFAULT now(),
CONSTRAINT product_status_chk CHECK (status IN ('active','merged','deprecated')),
CONSTRAINT product_quality_chk CHECK (quality_score >= 0 AND quality_score <= 1)
);
CREATE TABLE food_detail (
product_id UUID PRIMARY KEY REFERENCES product(id) ON DELETE CASCADE,
ingredients_text TEXT,
ingredients JSONB,
allergens TEXT[] NOT NULL DEFAULT '{}',
additives TEXT[] NOT NULL DEFAULT '{}',
nutriments JSONB,
nutrition_basis VARCHAR(16),
serving_size VARCHAR(32),
nutri_score CHAR(1),
labels TEXT[] NOT NULL DEFAULT '{}',
CONSTRAINT food_basis_chk CHECK (nutrition_basis IS NULL OR nutrition_basis IN ('per_100g','per_100ml','per_serving'))
);
-- Official manufacturer-suggested retail price snapshot (no purchase link).
CREATE TABLE product_msrp (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
product_id UUID NOT NULL REFERENCES product(id) ON DELETE CASCADE,
amount NUMERIC(12,2) NOT NULL,
currency CHAR(3) NOT NULL,
region VARCHAR(8) NOT NULL DEFAULT 'CN',
source_id UUID REFERENCES source(id),
source_url TEXT,
effective_date DATE,
note TEXT,
created_at TIMESTAMPTZ NOT NULL DEFAULT now()
);
CREATE TABLE product_image (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
product_id UUID NOT NULL REFERENCES product(id) ON DELETE CASCADE,
url TEXT NOT NULL,
kind VARCHAR(16) NOT NULL DEFAULT 'other',
license TEXT,
source_id UUID REFERENCES source(id),
CONSTRAINT image_kind_chk CHECK (kind IN ('front','ingredients','nutrition','other'))
);
-- Field-level provenance: which source provided which fields.
CREATE TABLE product_source (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
product_id UUID NOT NULL REFERENCES product(id) ON DELETE CASCADE,
source_id UUID REFERENCES source(id),
url TEXT,
fields TEXT[] NOT NULL DEFAULT '{}',
fetched_at TIMESTAMPTZ,
raw JSONB
);
CREATE TABLE merge_log (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
kept_id UUID,
merged_id UUID,
reason TEXT,
actor TEXT,
created_at TIMESTAMPTZ NOT NULL DEFAULT now()
);
-- Indexes
CREATE UNIQUE INDEX idx_product_gtin ON product (gtin) WHERE gtin IS NOT NULL;
CREATE INDEX idx_product_name_trgm ON product USING gin (name gin_trgm_ops);
CREATE INDEX idx_product_attrs ON product USING gin (attributes);
CREATE INDEX idx_product_tsv ON product USING gin (search_tsv);
CREATE INDEX idx_product_category ON product (category_id);
CREATE INDEX idx_product_brand ON product (brand_id);
CREATE INDEX idx_product_updated ON product (updated_at);
CREATE INDEX idx_food_nutriments ON food_detail USING gin (nutriments);
CREATE INDEX idx_category_path ON category USING gist (path);
CREATE INDEX idx_msrp_product ON product_msrp (product_id);
CREATE INDEX idx_psource_product ON product_source (product_id);
-- Keep search_tsv and updated_at in sync.
CREATE OR REPLACE FUNCTION product_sync_tsv() RETURNS trigger AS $$
BEGIN
NEW.search_tsv := to_tsvector('simple', coalesce(NEW.name, ''));
NEW.updated_at := now();
RETURN NEW;
END;
$$ LANGUAGE plpgsql;
CREATE TRIGGER trg_product_sync
BEFORE INSERT OR UPDATE ON product
FOR EACH ROW EXECUTE FUNCTION product_sync_tsv();
+2
View File
@@ -0,0 +1,2 @@
DELETE FROM attribute_definition;
DELETE FROM unit;
+29
View File
@@ -0,0 +1,29 @@
-- Unit dictionary seed. Keep factors aligned with ingestion/opengoods/units.py.
INSERT INTO unit (code, dimension, canonical, to_canonical_factor, aliases, display) VALUES
('mg', 'mass', 'g', 0.001, ARRAY['毫克'], 'mg'),
('g', 'mass', 'g', 1, ARRAY['','gram','grams'], 'g'),
('kg', 'mass', 'g', 1000, ARRAY['kgs','千克','公斤'], 'kg'),
('ml', 'volume', 'ml', 1, ARRAY['毫升','milliliter'], 'mL'),
('cl', 'volume', 'ml', 10, ARRAY['厘升'], 'cL'),
('l', 'volume', 'ml', 1000, ARRAY['L','','litre','liter'], 'L'),
('kj', 'energy', 'kJ', 1, ARRAY['kJ','千焦'], 'kJ'),
('kcal', 'energy', 'kJ', 4.184, ARRAY['千卡','大卡'], 'kcal'),
('pct', 'ratio', 'pct', 1, ARRAY['%','percent','百分比'], '%'),
('unit', 'count', 'unit',1, ARRAY['','','pcs','piece'], ''),
('mm', 'length', 'mm', 1, ARRAY['毫米'], 'mm'),
('cm', 'length', 'mm', 10, ARRAY['厘米'], 'cm'),
('day', 'duration', 'day', 1, ARRAY['','','days'], 'day')
ON CONFLICT (code) DO NOTHING;
-- A few common food attribute definitions referencing the unit dictionary.
INSERT INTO attribute_definition (key, label_zh, label_en, dimension, default_unit, aliases) VALUES
('energy', '能量', 'Energy', 'energy', 'kj', ARRAY['energy_kj']),
('proteins', '蛋白质', 'Proteins', 'mass', 'g', ARRAY['protein']),
('fat', '脂肪', 'Fat', 'mass', 'g', ARRAY['fats']),
('saturated_fat', '饱和脂肪','Saturated fat','mass', 'g', ARRAY['saturated-fat']),
('carbohydrates', '碳水化合物','Carbohydrates','mass', 'g', ARRAY['carbs']),
('sugars', '', 'Sugars', 'mass', 'g', ARRAY['sugar']),
('salt', '', 'Salt', 'mass', 'g', ARRAY['sodium_salt']),
('net_content', '净含量', 'Net content', NULL, NULL, ARRAY['quantity'])
ON CONFLICT (key) DO NOTHING;
+3
View File
@@ -0,0 +1,3 @@
-- remove seeded categories (children first via path depth)
DELETE FROM category_schema;
DELETE FROM category;
+53
View File
@@ -0,0 +1,53 @@
-- Seed a FOOD-focused category skeleton.
-- Structure = GS1 GPC backbone (segment/family/class) mapped to a self-built
-- Chinese tree. ltree labels are english slugs (ltree forbids spaces/CJK);
-- Chinese names live in name_zh. gpc_brick_code on leaves is a representative
-- starter value to be replaced by a full official GPC import later.
-- Root segment: Food/Beverage/Tobacco (GPC segment 50000000)
INSERT INTO category (name_zh, name_en, parent_id, path, gpc_brick_code, level)
VALUES ('食品饮料', 'Food/Beverage', NULL, 'food', '50000000', 0);
-- Families (level 1)
INSERT INTO category (name_zh, name_en, parent_id, path, gpc_brick_code, level)
SELECT v.name_zh, v.name_en, c.id, v.path::ltree, v.code, 1
FROM (VALUES
('饮料', 'Beverages', 'food.beverages', '50130000'),
('乳制品蛋类','Dairy/Eggs', 'food.dairy', '50180000'),
('烘焙', 'Bakery', 'food.bakery', '50100000'),
('零食', 'Snacks', 'food.snacks', '50190000'),
('粮油', 'Staples/Oils', 'food.staple', '50160000'),
('调味品', 'Condiments', 'food.condiments', '50170000')
) AS v(name_zh, name_en, path, code)
JOIN category c ON c.path = 'food';
-- Classes / leaves (level 2) with representative GPC brick codes
INSERT INTO category (name_zh, name_en, parent_id, path, gpc_brick_code, level)
SELECT v.name_zh, v.name_en, c.id, v.path::ltree, v.code, 2
FROM (VALUES
('包装饮用水', 'Bottled water', 'food.beverages.water', '10000224', 'food.beverages'),
('碳酸饮料', 'Carbonated', 'food.beverages.carbonated', '10000225', 'food.beverages'),
('果汁', 'Juice', 'food.beverages.juice', '10000226', 'food.beverages'),
('牛奶', 'Milk', 'food.dairy.milk', '10000158', 'food.dairy'),
('酸奶', 'Yogurt', 'food.dairy.yogurt', '10000159', 'food.dairy'),
('奶酪', 'Cheese', 'food.dairy.cheese', '10000160', 'food.dairy'),
('面包', 'Bread', 'food.bakery.bread', '10000040', 'food.bakery'),
('饼干', 'Biscuits', 'food.bakery.biscuits', '10000041', 'food.bakery'),
('薯片膨化', 'Chips/Snacks', 'food.snacks.chips', '10000310', 'food.snacks'),
('巧克力', 'Chocolate', 'food.snacks.chocolate', '10000311', 'food.snacks'),
('大米', 'Rice', 'food.staple.rice', '10000500', 'food.staple'),
('面条', 'Noodles', 'food.staple.noodles', '10000501', 'food.staple'),
('食用油', 'Cooking oil', 'food.staple.cooking_oil', '10000502', 'food.staple'),
('酱油', 'Soy sauce', 'food.condiments.soy_sauce', '10000600', 'food.condiments'),
('食盐', 'Table salt', 'food.condiments.salt', '10000601', 'food.condiments')
) AS v(name_zh, name_en, path, code, parent_path)
JOIN category c ON c.path = v.parent_path::ltree;
-- Parameter templates: leaf food categories use per_100g/ml nutrition basis.
INSERT INTO category_schema (category_id, required_attributes, recommended_attributes, nutriment_basis)
SELECT id,
ARRAY['net_content'],
ARRAY['energy','proteins','fat','carbohydrates','sugars','salt'],
CASE WHEN path <@ 'food.beverages' THEN 'per_100ml' ELSE 'per_100g' END
FROM category
WHERE level = 2;
+26 -1
View File
@@ -1,3 +1,28 @@
# Database migrations (golang-migrate) # Database migrations (golang-migrate)
SQL migrations live here from milestone M1. Format: `NNNN_description.up.sql` / `.down.sql`. SQL migrations for the OpenGoods database, applied with
[golang-migrate](https://github.com/golang-migrate/migrate).
Naming: `NNNN_description.up.sql` / `NNNN_description.down.sql`.
## Files
| Version | Up | 内容 |
|---------|----|------|
| 0001 | `0001_init` | 扩展(pgcrypto/pg_trgm/ltree) + 全部核心表 + 索引 + tsvector 触发器 |
| 0002 | `0002_seed_units` | 单位字典(与 `ingestion/opengoods/units.py` 一致)+ 常用营养参数定义 |
| 0003 | `0003_seed_categories` | 食品品类骨架(GS1 GPC 映射 + 自建中文树)+ 品类参数模板 |
## 运行
先起本地依赖:`docker compose up -d postgres`
```bash
export DBURL="postgres://opengoods:opengoods@localhost:5432/opengoods?sslmode=disable"
migrate -path migrations -database "$DBURL" up # 升级到最新
migrate -path migrations -database "$DBURL" down -all # 全部回滚
migrate -path migrations -database "$DBURL" version # 查看当前版本
```
安装 CLI`go install -tags 'postgres' github.com/golang-migrate/migrate/v4/cmd/migrate@v4.18.1`
> ltree 标签为英文 slug(不支持空格/中文),中文名存于 `category.name_zh`。
> `gpc_brick_code` 为食品子集的代表值,后续用官方 GPC 全量导入替换。