Download src/knowledge_parsing/parse.py from DataEyond/Agentic-Service-Data-Eyond-Catalog: direct link, hf CLI and curl.
- Browser
- Download file 14.3 kB
-
https://huggingface.co/spaces/DataEyond/Agentic-Service-Data-Eyond-Catalog/resolve/main/src/knowledge_parsing/parse.py
- Command line
-
hf download hf://spaces/DataEyond/Agentic-Service-Data-Eyond-Catalog/src/knowledge_parsing/parse.py
-
curl -L -o parse.py https://huggingface.co/spaces/DataEyond/Agentic-Service-Data-Eyond-Catalog/resolve/main/src/knowledge_parsing/parse.py
14.3 kB
| """Stage 2 — parse a document with MinerU, cached by content. | |
| Why the cache is keyed to document CONTENT rather than to a run id: parsing is | |
| by far the most expensive stage (19 s/page on CPU). A per-run cache would | |
| re-parse everything each time the normaliser changes and the pipeline is run | |
| again, for nothing. Keyed by content + settings + version, the same document is | |
| parsed once — however many times the normaliser is revised. | |
| Cache key = sha256(file bytes) + settings fingerprint + MinerU version. The | |
| version is in there so figures from an older MinerU never quietly mix with | |
| figures from a newer one — which matters, because those figures are used to | |
| justify spend. | |
| """ | |
| from __future__ import annotations | |
| import hashlib | |
| import json | |
| import shutil | |
| import time | |
| from dataclasses import dataclass | |
| from pathlib import Path | |
| from .config import PipelineConfig | |
| class ParseResult: | |
| doc_id: str | |
| source: Path | |
| cache_dir: Path # folder holding MinerU's output untouched | |
| content_list: Path # *_content_list.json | |
| middle_json: Path | None | |
| pages: int | |
| seconds: float | |
| from_cache: bool | |
| backend_recorded: str | None # read from _middle.json, NOT from the config | |
| mineru_version: str | None | |
| def mineru_version() -> str: | |
| try: | |
| from mineru.version import __version__ # type: ignore | |
| return str(__version__) | |
| except Exception: | |
| try: | |
| from importlib.metadata import version | |
| return version("mineru") | |
| except Exception: | |
| return "unknown" | |
| def file_hash(path: Path, length: int = 16) -> str: | |
| h = hashlib.sha256() | |
| with path.open("rb") as f: | |
| for block in iter(lambda: f.read(1 << 20), b""): | |
| h.update(block) | |
| return h.hexdigest()[:length] | |
| def cache_key(source: Path, cfg: PipelineConfig) -> str: | |
| # The parser version goes in the key so figures from two parser versions never | |
| # mix. For Mistral that version is the model id, not the MinerU build. | |
| parser_ver = ( | |
| cfg.mistral_model if cfg.backend == "mistral" | |
| else cfg.paddle_model if cfg.backend == "paddleocr" | |
| else mineru_version() | |
| ) | |
| return f"{file_hash(source)}-{cfg.fingerprint()}-{parser_ver}" | |
| def _mineru_name(doc_id: str, source: Path) -> str: | |
| """A length-bounded folder name for MinerU's own output. | |
| MinerU names its output folder after whatever is passed as | |
| `pdf_file_names`, and then writes images into it under a full SHA-256 | |
| filename. Passing the real document name blows the Windows MAX_PATH limit | |
| of 260: measured on this repo, the fixed part of the path costs 217 | |
| characters, so the name may not exceed 43. `STD_2026_006_MNO_ Rev.0.0 - | |
| Production Parameter and ECA` is 56 and fails. | |
| **Cosmetic only.** Uniqueness is already guaranteed by the parent folder, | |
| which is `cache_key` — the full SHA-256 of the file contents. Two different | |
| documents can never share a cache folder however alike their names. The | |
| readable part exists so the cache can be browsed while debugging, nothing | |
| more. | |
| Head AND tail are kept, because the names that actually collide share a | |
| long prefix and differ at the end — "… and ECA" versus "… and ECA (2)". | |
| Keeping only the head would render exactly those indistinguishable. | |
| **Temporary.** Once `doc_id` becomes the content hash rather than the file | |
| name (see the README), every path here is short by construction and this | |
| function goes away. It is deliberately kept simple for that reason. | |
| """ | |
| safe = "".join(c if c.isalnum() or c in " ._-()" else "_" for c in doc_id).strip() | |
| short = safe if len(safe) <= 20 else f"{safe[:12]}~{safe[-8:]}" | |
| return f"{short}-{file_hash(source, 8)}" | |
| def require_gpu(cfg: PipelineConfig) -> None: | |
| """`hybrid` and `vlm` need a GPU — checked only when MinerU is about to run. | |
| Deliberately not a config-time check. A cache hit needs no GPU, and | |
| rebuilding artifacts from MinerU's saved output is the normal way to pick up | |
| a normalisation fix without paying for a re-parse: the expensive stage | |
| (PDF -> content_list) is cached, and only the cheap CPU stage re-runs. | |
| Refusing at startup would block exactly that. | |
| """ | |
| # RuntimeError, not SystemExit: the batch may hold a mix of cached and | |
| # uncached documents, and one that needs a GPU must not abort the ones that | |
| # do not. It is recorded as a per-document failure like any other. | |
| if cfg.backend in {"hybrid", "vlm"} and shutil.which("nvidia-smi") is None: | |
| raise RuntimeError( | |
| f"Backend '{cfg.backend}' needs a GPU to parse, and nvidia-smi was not found.\n" | |
| " -> This document is not in the parse cache, so MinerU has to run.\n" | |
| " -> To parse on a laptop without one, ask for pipeline:\n" | |
| " python -m src.knowledge_parsing.run " | |
| "--input data/knowledge_docs/ --backend pipeline\n" | |
| " -> That output is valid, but has NO image descriptions and the\n" | |
| " multiplication sign is copied as the letter 'x'. Do not use\n" | |
| " it as a reference measurement." | |
| ) | |
| def _find_output(root: Path) -> tuple[Path | None, Path | None]: | |
| """Locate content_list & middle json anywhere under root. | |
| Deliberately a glob rather than a fixed path: MinerU's subfolder layout | |
| differs per backend ('auto' for pipeline, something else for vlm). Globbing | |
| means this module needs no change when the backend does. | |
| """ | |
| content = next(iter(sorted(root.rglob("*_content_list.json"))), None) | |
| middle = next(iter(sorted(root.rglob("*_middle.json"))), None) | |
| return content, middle | |
| def _read_middle(middle: Path | None) -> tuple[str | None, str | None]: | |
| """The backend & version that ACTUALLY ran, taken from MinerU's own output.""" | |
| if not middle or not middle.exists(): | |
| return None, None | |
| try: | |
| d = json.loads(middle.read_text(encoding="utf-8")) | |
| return d.get("_backend"), d.get("_version_name") | |
| except Exception: | |
| return None, None | |
| def _parse_mistral( | |
| source: Path, cfg: PipelineConfig, doc_id: str, target: Path, marker: Path | |
| ) -> ParseResult: | |
| """Parse one document with Mistral OCR, cached exactly like the MinerU path. | |
| The cached output is a single `mistral_ocr.json` (the API response) rather | |
| than MinerU's `*_content_list.json`; `ParseResult.content_list` points at it | |
| and run.py dispatches on `cfg.backend` to the matching normaliser. | |
| """ | |
| out = target / "mistral_ocr.json" | |
| if marker.exists() and out.exists(): | |
| meta = json.loads(marker.read_text(encoding="utf-8")) | |
| return ParseResult( | |
| doc_id=doc_id, source=source, cache_dir=target, | |
| content_list=out, middle_json=None, | |
| pages=meta.get("pages", 0), seconds=meta.get("seconds", 0.0), | |
| from_cache=True, backend_recorded="mistral-ocr", | |
| mineru_version=meta.get("model") or cfg.mistral_model, | |
| ) | |
| from .mistral_ocr import run_ocr | |
| started = time.perf_counter() | |
| response = run_ocr(source, cfg) # no GPU; needs MISTRAL_API_KEY | |
| seconds = time.perf_counter() - started | |
| pages = len(response.get("pages", [])) | |
| model = response.get("model") or cfg.mistral_model | |
| staging = target.with_name(target.name + ".partial") | |
| shutil.rmtree(staging, ignore_errors=True) | |
| staging.mkdir(parents=True, exist_ok=True) | |
| (staging / "mistral_ocr.json").write_text( | |
| json.dumps(response, ensure_ascii=False, indent=2), encoding="utf-8" | |
| ) | |
| staging.rename(target) | |
| marker.write_text( | |
| json.dumps({"pages": pages, "seconds": seconds, "doc_id": doc_id, "model": model}), | |
| encoding="utf-8", | |
| ) | |
| return ParseResult( | |
| doc_id=doc_id, source=source, cache_dir=target, | |
| content_list=target / "mistral_ocr.json", middle_json=None, | |
| pages=pages, seconds=seconds, from_cache=False, | |
| backend_recorded="mistral-ocr", mineru_version=model, | |
| ) | |
| def _parse_paddle( | |
| source: Path, cfg: PipelineConfig, doc_id: str, target: Path, marker: Path | |
| ) -> ParseResult: | |
| """Parse one document with PaddleOCR-VL, cached exactly like the Mistral path. | |
| The cached `paddle_ocr.json` is `paddle_ocr.run_ocr`'s output — the page results | |
| with every figure's bytes already embedded, so a cache hit needs no network. | |
| """ | |
| out = target / "paddle_ocr.json" | |
| if marker.exists() and out.exists(): | |
| meta = json.loads(marker.read_text(encoding="utf-8")) | |
| return ParseResult( | |
| doc_id=doc_id, source=source, cache_dir=target, | |
| content_list=out, middle_json=None, | |
| pages=meta.get("pages", 0), seconds=meta.get("seconds", 0.0), | |
| from_cache=True, backend_recorded="paddleocr-vl", | |
| mineru_version=meta.get("model") or cfg.paddle_model, | |
| ) | |
| from .paddle_ocr import run_ocr | |
| started = time.perf_counter() | |
| response = run_ocr(source, cfg) # hosted API; needs PADDLEOCR_ACCESS_TOKEN | |
| seconds = time.perf_counter() - started | |
| pages = len(response.get("pages", [])) | |
| model = response.get("model") or cfg.paddle_model | |
| staging = target.with_name(target.name + ".partial") | |
| shutil.rmtree(staging, ignore_errors=True) | |
| staging.mkdir(parents=True, exist_ok=True) | |
| (staging / "paddle_ocr.json").write_text( | |
| json.dumps(response, ensure_ascii=False), encoding="utf-8" | |
| ) | |
| shutil.rmtree(target, ignore_errors=True) | |
| staging.rename(target) | |
| marker.write_text( | |
| json.dumps({"pages": pages, "seconds": seconds, "doc_id": doc_id, "model": model}), | |
| encoding="utf-8", | |
| ) | |
| return ParseResult( | |
| doc_id=doc_id, source=source, cache_dir=target, | |
| content_list=out, middle_json=None, | |
| pages=pages, seconds=seconds, from_cache=False, | |
| backend_recorded="paddleocr-vl", mineru_version=model, | |
| ) | |
| def parse_document(source: Path, cfg: PipelineConfig) -> ParseResult: | |
| """Parse one document. A cache hit is not re-run.""" | |
| doc_id = source.stem | |
| target = cfg.cache_dir / "parse" / cache_key(source, cfg) | |
| marker = target / ".done" | |
| if cfg.backend == "mistral": | |
| return _parse_mistral(source, cfg, doc_id, target, marker) | |
| if cfg.backend == "paddleocr": | |
| return _parse_paddle(source, cfg, doc_id, target, marker) | |
| if marker.exists(): | |
| content, middle = _find_output(target) | |
| if content: | |
| meta = json.loads(marker.read_text(encoding="utf-8")) | |
| backend_recorded, version = _read_middle(middle) | |
| return ParseResult( | |
| doc_id=doc_id, source=source, cache_dir=target, | |
| content_list=content, middle_json=middle, | |
| pages=meta.get("pages", 0), seconds=meta.get("seconds", 0.0), | |
| from_cache=True, | |
| backend_recorded=backend_recorded, mineru_version=version, | |
| ) | |
| shutil.rmtree(target, ignore_errors=True) # cache corrupt, redo it | |
| from mineru.cli.common import do_parse, read_fn | |
| # MinerU dispatches on prefixes: "pipeline", then `startswith("vlm-")` and | |
| # `startswith("hybrid-")`. The bare names "vlm"/"hybrid" match NONE of those | |
| # branches — MinerU simply does nothing, without error, and the failure only | |
| # surfaces here as "content_list not found", which points at the wrong place | |
| # entirely. The short names stay in the config because they are part of the | |
| # cache key and read better; the translation happens here, right before | |
| # MinerU is touched. | |
| mineru_backend = { | |
| "pipeline": "pipeline", | |
| "vlm": "vlm-engine", | |
| "hybrid": "hybrid-engine", | |
| }[cfg.backend] | |
| # The playground does this before parsing; aliases like "latin"/"en" are | |
| # resolved to a canonical code. Matched here so results are identical. | |
| # Cache missed, so MinerU really is about to run — this is the point where a | |
| # GPU stops being optional. | |
| require_gpu(cfg) | |
| lang = cfg.canonical_lang() | |
| # with_suffix() REPLACES the last suffix, and the folder name ends in the | |
| # MinerU version ("...-3.4.4") — so ".4" would be thrown away and two | |
| # different versions would collide on one temp directory. | |
| staging = target.with_name(target.name + ".partial") | |
| shutil.rmtree(staging, ignore_errors=True) | |
| staging.mkdir(parents=True, exist_ok=True) | |
| pdf_bytes = read_fn(source) | |
| started = time.perf_counter() | |
| do_parse( | |
| output_dir=str(staging), | |
| # NOT doc_id: this becomes a folder name, and a long document name | |
| # exceeds MAX_PATH once MinerU appends a SHA-256 image filename to it. | |
| pdf_file_names=[_mineru_name(doc_id, source)], | |
| pdf_bytes_list=[pdf_bytes], | |
| p_lang_list=[lang], | |
| backend=mineru_backend, | |
| formula_enable=cfg.formula_enable, | |
| table_enable=cfg.table_enable, | |
| f_draw_layout_bbox=cfg.write_debug_pdf, | |
| f_draw_span_bbox=cfg.write_debug_pdf, | |
| start_page_id=cfg.start_page, | |
| end_page_id=cfg.end_page, | |
| **({"effort": cfg.effort} if cfg.backend == "hybrid" else {}), | |
| ) | |
| seconds = time.perf_counter() - started | |
| content, middle = _find_output(staging) | |
| if content is None: | |
| raise RuntimeError( | |
| f"MinerU finished but no *_content_list.json was found in {staging}" | |
| ) | |
| pages = _count_pages(source, cfg) | |
| staging.rename(target) | |
| content = target / content.relative_to(staging) | |
| middle = target / middle.relative_to(staging) if middle else None | |
| marker.write_text( | |
| json.dumps({"pages": pages, "seconds": seconds, "doc_id": doc_id}), | |
| encoding="utf-8", | |
| ) | |
| backend_recorded, version = _read_middle(middle) | |
| return ParseResult( | |
| doc_id=doc_id, source=source, cache_dir=target, | |
| content_list=content, middle_json=middle, | |
| pages=pages, seconds=seconds, from_cache=False, | |
| backend_recorded=backend_recorded, mineru_version=version, | |
| ) | |
| def _count_pages(source: Path, cfg: PipelineConfig) -> int: | |
| if cfg.end_page is not None: | |
| return cfg.end_page - cfg.start_page + 1 | |
| try: | |
| from pypdf import PdfReader | |
| return len(PdfReader(str(source)).pages) | |
| except Exception: | |
| return 0 | |