From 8dea386bc7265fffe22b898021841af1877f63fb Mon Sep 17 00:00:00 2001 From: Alikhan Serik Date: Tue, 14 Oct 2025 15:16:56 +0500 Subject: [PATCH 1/4] Add gitignore .idea/ for Pycharm --- .gitignore | 3 +++ 1 file changed, 3 insertions(+) diff --git a/.gitignore b/.gitignore index 9a5aced..95d7b25 100644 --- a/.gitignore +++ b/.gitignore @@ -137,3 +137,6 @@ dist # Vite logs files vite.config.js.timestamp-* vite.config.ts.timestamp-* + +# PyCharm +.idea/ \ No newline at end of file From d0a4a1556de58751b7e8ecd66c9eb0851cdef094 Mon Sep 17 00:00:00 2001 From: Alikhan Serik Date: Sat, 29 Nov 2025 15:03:28 +0500 Subject: [PATCH 2/4] feat(parser): working response from API --- .../openfoodfacts_parser.py | 19 +++++++++++++++++++ 1 file changed, 19 insertions(+) create mode 100644 domains/world.openfoodfacts.org/another_implementation/openfoodfacts_parser.py diff --git a/domains/world.openfoodfacts.org/another_implementation/openfoodfacts_parser.py b/domains/world.openfoodfacts.org/another_implementation/openfoodfacts_parser.py new file mode 100644 index 0000000..af97c0e --- /dev/null +++ b/domains/world.openfoodfacts.org/another_implementation/openfoodfacts_parser.py @@ -0,0 +1,19 @@ +import httpx +import asyncio +import json + +async def main(): + url = "https://world.openfoodfacts.org/api/v2/search" + + params = { + "page": 1, + "page_size": 10, + "fields": "code,product_name,image_url,created_t,last_modified_t,last_updated_t", + } + + async with httpx.AsyncClient(timeout=20.0) as client: + response = await client.get(url, params=params) + print("STATUS:", response.status_code) + print(json.dumps(response.json(), indent=2)) + +asyncio.run(main()) From e93b194d850828c8bf8f84cb2b3e4271ab796c43 Mon Sep 17 00:00:00 2001 From: Alikhan Serik Date: Sat, 29 Nov 2025 19:02:45 +0500 Subject: [PATCH 3/4] feat(parser): working parser --- .../openfoodfacts_parser.py | 163 ++++++++++++++++-- 1 file changed, 152 insertions(+), 11 deletions(-) diff --git a/domains/world.openfoodfacts.org/another_implementation/openfoodfacts_parser.py b/domains/world.openfoodfacts.org/another_implementation/openfoodfacts_parser.py index af97c0e..8c0961f 100644 --- a/domains/world.openfoodfacts.org/another_implementation/openfoodfacts_parser.py +++ b/domains/world.openfoodfacts.org/another_implementation/openfoodfacts_parser.py @@ -1,19 +1,160 @@ -import httpx import asyncio +import httpx import json +import os +from datetime import datetime -async def main(): - url = "https://world.openfoodfacts.org/api/v2/search" +STATE_FILE = "state.json" +OUTPUT_FILE = "products.json" + +API_URL = "https://world.openfoodfacts.org/api/v2/search" + +CONCURRENCY = 10 # параллельные запросы +PAGE_SIZE = 100 # чем больше — тем быстрее, но риск ошибок выше +STOP_IF_EMPTY = True # если API вернул пустую страницу — стоп + + +# ----------------------------- +# STATE MANAGER +# ----------------------------- +def load_state(): + if not os.path.exists(STATE_FILE): + return {"last_updated_t": 0, "page": 1} + + with open(STATE_FILE, "r") as f: + return json.load(f) + + +def save_state(state): + with open(STATE_FILE, "w") as f: + json.dump(state, f, indent=2) + + +# ----------------------------- +# JSON STREAM WRITER +# ----------------------------- +def ensure_output_file(): + if not os.path.exists(OUTPUT_FILE): + with open(OUTPUT_FILE, "w") as f: + f.write("[\n") + +def append_product(product): + with open(OUTPUT_FILE, "a", encoding="utf-8") as f: + json.dump(product, f, ensure_ascii=False) + f.write(",\n") + + +# ----------------------------- +# API REQUEST +# ----------------------------- +async def fetch_page(client: httpx.AsyncClient, page: int, last_updated_t: int): params = { - "page": 1, - "page_size": 10, - "fields": "code,product_name,image_url,created_t,last_modified_t,last_updated_t", + "page": page, + "page_size": PAGE_SIZE, + "fields": "code,product_name,image_url,last_updated_t", + "sort_by": "last_updated_t", + "sort_order": "asc", + "last_updated_t": f">{last_updated_t}" } - async with httpx.AsyncClient(timeout=20.0) as client: - response = await client.get(url, params=params) - print("STATUS:", response.status_code) - print(json.dumps(response.json(), indent=2)) + r = await client.get(API_URL, params=params) + + if r.status_code != 200: + print(f"[ERROR] page {page} status {r.status_code}") + return None + + return r.json() + + +# ----------------------------- +# EXTRACT PRODUCT +# ----------------------------- +def extract_product(p): + return { + "barcode": p.get("code"), + "name": p.get("product_name"), + "image_links": [p.get("image_url")], + "updated_at": datetime.utcfromtimestamp(p["last_updated_t"]).isoformat() + } + + +# ----------------------------- +# WORKER +# ----------------------------- +async def worker(page_queue: asyncio.Queue, state): + async with httpx.AsyncClient(timeout=30.0) as client: + while True: + page = await page_queue.get() + if page is None: + return + + data = await fetch_page(client, page, state["last_updated_t"]) + + if not data or "products" not in data: + page_queue.task_done() + continue + + products = data["products"] + + if STOP_IF_EMPTY and len(products) == 0: + print("Reached empty page → stopping early.") + page_queue.task_done() + return + + latest_ts = state["last_updated_t"] + + for p in products: + if "last_updated_t" not in p: + continue + + ts = p["last_updated_t"] + if ts > latest_ts: + latest_ts = ts + + append_product(extract_product(p)) + + state["last_updated_t"] = latest_ts + state["page"] = page + save_state(state) + + print(f"[PAGE {page}] processed {len(products)} products") + + page_queue.task_done() + + +# ----------------------------- +# MAIN LOOP +# ----------------------------- +async def main(): + state = load_state() + ensure_output_file() + + start_page = state["page"] + print(f"Resuming from page {start_page}, last_updated_t={state['last_updated_t']}") + + page_queue = asyncio.Queue() + + for page in range(start_page, 10_000_000): + page_queue.put_nowait(page) + + tasks = [] + for _ in range(CONCURRENCY): + t = asyncio.create_task(worker(page_queue, state)) + tasks.append(t) + + await page_queue.join() + + for _ in range(CONCURRENCY): + page_queue.put_nowait(None) + + await asyncio.gather(*tasks) + + print("DONE. Close JSON manually with ] when finished parsing.") + -asyncio.run(main()) +if __name__ == "__main__": + try: + asyncio.run(main()) + except KeyboardInterrupt: + print("Stopped manually.") From be4351929e0ca47084957b9fbfa0639c22e22286 Mon Sep 17 00:00:00 2001 From: Alikhan Serik Date: Mon, 1 Dec 2025 19:58:41 +0500 Subject: [PATCH 4/4] feat(parser): OpenFoodFacts parser without testing --- .gitignore | 5 +- .../openfoodfacts_parser.py | 141 +++++++++--------- .../another_implementation/state.json | 3 + 3 files changed, 78 insertions(+), 71 deletions(-) create mode 100644 domains/world.openfoodfacts.org/another_implementation/state.json diff --git a/.gitignore b/.gitignore index 95d7b25..7ec9e8d 100644 --- a/.gitignore +++ b/.gitignore @@ -139,4 +139,7 @@ vite.config.js.timestamp-* vite.config.ts.timestamp-* # PyCharm -.idea/ \ No newline at end of file +.idea/ + +# Large files +domains/world.openfoodfacts.org/another_implementation/products.json \ No newline at end of file diff --git a/domains/world.openfoodfacts.org/another_implementation/openfoodfacts_parser.py b/domains/world.openfoodfacts.org/another_implementation/openfoodfacts_parser.py index 8c0961f..991f196 100644 --- a/domains/world.openfoodfacts.org/another_implementation/openfoodfacts_parser.py +++ b/domains/world.openfoodfacts.org/another_implementation/openfoodfacts_parser.py @@ -9,17 +9,17 @@ API_URL = "https://world.openfoodfacts.org/api/v2/search" -CONCURRENCY = 10 # параллельные запросы -PAGE_SIZE = 100 # чем больше — тем быстрее, но риск ошибок выше -STOP_IF_EMPTY = True # если API вернул пустую страницу — стоп +CONCURRENCY = 20 +PAGE_SIZE = 100 +MAX_PAGES_PER_CYCLE = 200 # ----------------------------- -# STATE MANAGER +# STATE # ----------------------------- def load_state(): if not os.path.exists(STATE_FILE): - return {"last_updated_t": 0, "page": 1} + return {"last_updated_t": 0} with open(STATE_FILE, "r") as f: return json.load(f) @@ -31,7 +31,7 @@ def save_state(state): # ----------------------------- -# JSON STREAM WRITER +# OUTPUT # ----------------------------- def ensure_output_file(): if not os.path.exists(OUTPUT_FILE): @@ -39,16 +39,16 @@ def ensure_output_file(): f.write("[\n") -def append_product(product): +def append_product(p): with open(OUTPUT_FILE, "a", encoding="utf-8") as f: - json.dump(product, f, ensure_ascii=False) + json.dump(p, f, ensure_ascii=False) f.write(",\n") # ----------------------------- -# API REQUEST +# API REQUEST # ----------------------------- -async def fetch_page(client: httpx.AsyncClient, page: int, last_updated_t: int): +async def fetch_page(client: httpx.AsyncClient, page: int, last_updated_t: int, retries=3): params = { "page": page, "page_size": PAGE_SIZE, @@ -58,17 +58,23 @@ async def fetch_page(client: httpx.AsyncClient, page: int, last_updated_t: int): "last_updated_t": f">{last_updated_t}" } - r = await client.get(API_URL, params=params) + for attempt in range(retries): + try: + r = await client.get(API_URL, params=params) + if r.status_code == 200: + return r.json() + return None - if r.status_code != 200: - print(f"[ERROR] page {page} status {r.status_code}") - return None + except (httpx.ReadTimeout, httpx.ConnectTimeout): + print(f"[TIMEOUT] page {page}, attempt {attempt+1}/{retries}") + await asyncio.sleep(1 + attempt) - return r.json() + print(f"[FAILED] page {page}") + return None # ----------------------------- -# EXTRACT PRODUCT +# PRODUCT CLEANER # ----------------------------- def extract_product(p): return { @@ -80,81 +86,76 @@ def extract_product(p): # ----------------------------- -# WORKER +# FETCH MANY PAGES IN PARALLEL # ----------------------------- -async def worker(page_queue: asyncio.Queue, state): - async with httpx.AsyncClient(timeout=30.0) as client: - while True: - page = await page_queue.get() - if page is None: - return +async def fetch_cycle(last_updated_t): + async with httpx.AsyncClient(timeout=60.0) as client: + sem = asyncio.Semaphore(CONCURRENCY) - data = await fetch_page(client, page, state["last_updated_t"]) + async def wrapped(page): + async with sem: + return await fetch_page(client, page, last_updated_t) + tasks = [asyncio.create_task(wrapped(page)) for page in range(1, MAX_PAGES_PER_CYCLE + 1)] + return await asyncio.gather(*tasks) + + +# ----------------------------- +# MAIN LOOP +# ----------------------------- +async def main(): + state = load_state() + ensure_output_file() + + print(f"Starting with last_updated_t={state['last_updated_t']}") + + while True: + batch = await fetch_cycle(state["last_updated_t"]) + + all_products = [] + max_ts = state["last_updated_t"] + empty_pages = 0 + + for idx, data in enumerate(batch, start=1): if not data or "products" not in data: - page_queue.task_done() + empty_pages += 1 continue products = data["products"] - if STOP_IF_EMPTY and len(products) == 0: - print("Reached empty page → stopping early.") - page_queue.task_done() - return + if len(products) == 0: + empty_pages += 1 + continue - latest_ts = state["last_updated_t"] + print(f"[PAGE {idx}] +{len(products)} products") for p in products: if "last_updated_t" not in p: continue ts = p["last_updated_t"] - if ts > latest_ts: - latest_ts = ts - - append_product(extract_product(p)) - - state["last_updated_t"] = latest_ts - state["page"] = page - save_state(state) - - print(f"[PAGE {page}] processed {len(products)} products") - - page_queue.task_done() - - -# ----------------------------- -# MAIN LOOP -# ----------------------------- -async def main(): - state = load_state() - ensure_output_file() - - start_page = state["page"] - print(f"Resuming from page {start_page}, last_updated_t={state['last_updated_t']}") - - page_queue = asyncio.Queue() + if ts > max_ts: + max_ts = ts - for page in range(start_page, 10_000_000): - page_queue.put_nowait(page) + all_products.append(extract_product(p)) - tasks = [] - for _ in range(CONCURRENCY): - t = asyncio.create_task(worker(page_queue, state)) - tasks.append(t) + # no new data -> stop + if empty_pages == MAX_PAGES_PER_CYCLE: + print("No more new data. Stopping.") + break - await page_queue.join() + # write downloaded products + for pr in all_products: + append_product(pr) - for _ in range(CONCURRENCY): - page_queue.put_nowait(None) + print(f"Downloaded {len(all_products)} new products in this cycle.") + print(f"Updated last_updated_t → {max_ts}") - await asyncio.gather(*tasks) + state["last_updated_t"] = max_ts + save_state(state) - print("DONE. Close JSON manually with ] when finished parsing.") + print("Parsing finished. Don't forget to close array with ]") if __name__ == "__main__": - try: - asyncio.run(main()) - except KeyboardInterrupt: - print("Stopped manually.") + asyncio.run(main()) diff --git a/domains/world.openfoodfacts.org/another_implementation/state.json b/domains/world.openfoodfacts.org/another_implementation/state.json new file mode 100644 index 0000000..7c50f59 --- /dev/null +++ b/domains/world.openfoodfacts.org/another_implementation/state.json @@ -0,0 +1,3 @@ +{ + "last_updated_t": 1764492189 +} \ No newline at end of file