Coverage for oc_meta / run / migration / rdf_to_nquads.py: 100%
30 statements
« prev ^ index » next coverage.py v7.13.4, created at 2026-07-25 10:39 +0000
« prev ^ index » next coverage.py v7.13.4, created at 2026-07-25 10:39 +0000
1#!/usr/bin/env python
3# Copyright 2026 Arcangelo Massari <arcangelo.massari@unibo.it>
4# SPDX-FileCopyrightText: 2026 Arcangelo Massari <arcangelo.massari@unibo.it>
5#
6# SPDX-License-Identifier: ISC
8import argparse
9import multiprocessing
10import os
11import zipfile
12from concurrent.futures import ProcessPoolExecutor
13from functools import partial
14from pathlib import Path
16import py7zr
17from rdflib import Dataset
18from rich_argparse import RichHelpFormatter
20from oc_meta.lib.console import console, create_progress
21from oc_meta.lib.file_manager import collect_zip_files
24def process_zip_file(
25 zip_path: Path, output_dir: Path, input_dir_path: Path, compress: bool
26) -> None:
27 graph = Dataset(default_union=True)
28 with zipfile.ZipFile(zip_path, "r") as zf:
29 json_file = next(name for name in zf.namelist() if name.endswith(".json"))
30 with zf.open(json_file) as f:
31 graph.parse(f, format="json-ld")
33 nquads_output = graph.serialize(format="nquads")
35 relative_path = zip_path.relative_to(input_dir_path)
36 output_filename = str(relative_path).replace(os.sep, "-")
37 output_filename = Path(output_filename).with_suffix(".nq").name
38 output_nq_path = output_dir / output_filename
40 with open(output_nq_path, "w", encoding="utf-8") as f:
41 f.write(nquads_output)
43 if compress:
44 output_7z_path = output_nq_path.with_suffix(".nq.7z")
45 with py7zr.SevenZipFile(output_7z_path, "w") as archive:
46 archive.write(output_nq_path, output_filename)
47 output_nq_path.unlink()
50def main() -> None: # pragma: no cover
51 parser = argparse.ArgumentParser(
52 description="Converts JSON-LD files from ZIP archives to N-Quads format.",
53 formatter_class=RichHelpFormatter,
54 )
55 parser.add_argument(
56 "input_dir",
57 type=str,
58 help="Input directory containing ZIP files (recursive search)",
59 )
60 parser.add_argument(
61 "output_dir", type=str, help="Output directory for the converted .nq files"
62 )
63 parser.add_argument(
64 "-m",
65 "--mode",
66 type=str,
67 choices=["all", "data", "prov"],
68 default="all",
69 help="Mode: 'all' for all ZIP files (default), 'data' for entity data only, 'prov' for provenance only",
70 )
71 parser.add_argument(
72 "-w",
73 "--workers",
74 type=int,
75 default=None,
76 help="Number of worker processes (defaults to CPU count)",
77 )
78 parser.add_argument(
79 "-c",
80 "--compress",
81 action="store_true",
82 help="Compress output files using 7z format",
83 )
84 args = parser.parse_args()
86 input_path = Path(args.input_dir).resolve()
87 output_path = Path(args.output_dir).resolve()
88 num_workers = args.workers if args.workers else multiprocessing.cpu_count()
90 output_path.mkdir(parents=True, exist_ok=True)
92 zip_files = collect_zip_files(
93 str(input_path),
94 only_data=args.mode == "data",
95 only_prov=args.mode == "prov",
96 )
97 total_files = len(zip_files)
99 mode_labels = {"all": "", "data": "data ", "prov": "provenance "}
100 console.print(
101 f"Found {total_files} {mode_labels[args.mode]}ZIP files in {input_path}"
102 )
103 console.print(f"Output directory: {output_path}")
104 console.print(f"Workers: {num_workers}")
105 console.print(f"Compression: {'7z' if args.compress else 'none'}")
107 fail_count = 0
108 task_func = partial(
109 process_zip_file,
110 output_dir=output_path,
111 input_dir_path=input_path,
112 compress=args.compress,
113 )
115 # Use forkserver to avoid deadlocks when forking in a multi-threaded environment
116 ctx = multiprocessing.get_context("forkserver")
117 with ProcessPoolExecutor(max_workers=num_workers, mp_context=ctx) as executor:
118 iterator = executor.map(task_func, zip_files)
120 with create_progress() as progress:
121 task = progress.add_task("Converting", total=total_files)
122 while True:
123 try:
124 next(iterator)
125 progress.update(task, advance=1)
126 except StopIteration:
127 break
128 except Exception as e:
129 console.print(f"[red]Error: {e}[/red]")
130 fail_count += 1
131 progress.update(task, advance=1)
133 console.print()
134 console.print("Final report")
135 console.print(f" Success: {total_files - fail_count}")
136 console.print(f" Failed: {fail_count}")
139if __name__ == "__main__": # pragma: no cover
140 main()