Coverage for oc_meta / run / migration / stream_nquads.py: 95%

74 statements  

« prev     ^ index     » next       coverage.py v7.13.4, created at 2026-07-25 10:39 +0000

1#!/usr/bin/env python 

2 

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 

7 

8import argparse 

9import gzip 

10import multiprocessing 

11import sys 

12import zipfile 

13from collections.abc import Iterable, Iterator 

14from pathlib import Path 

15from tempfile import TemporaryDirectory 

16from typing import IO 

17 

18from rdflib import Dataset 

19from rich.console import Console 

20from rich.progress import ( 

21 BarColumn, 

22 MofNCompleteColumn, 

23 Progress, 

24 SpinnerColumn, 

25 TaskID, 

26 TextColumn, 

27 TimeElapsedColumn, 

28 TimeRemainingColumn, 

29) 

30from rich_argparse import RichHelpFormatter 

31 

32from oc_meta.lib.file_manager import collect_zip_files 

33 

34DEFAULT_LINES_PER_FILE = 10000000 

35 

36 

37def convert_zip_to_nquads(zip_path: str) -> bytes: 

38 try: 

39 with zipfile.ZipFile(zip_path, "r") as zf: 

40 json_file = next(n for n in zf.namelist() if n.endswith(".json")) 

41 graph = Dataset(default_union=True) 

42 with zf.open(json_file) as f: 

43 graph.parse(f, format="json-ld") 

44 return graph.serialize(format="nquads").encode("utf-8") 

45 except Exception: 

46 print(f"Failed to convert: {zip_path}", file=sys.stderr, flush=True) 

47 raise 

48 

49 

50def convert_zip_to_nquads_file(task: tuple[str, str]) -> str: 

51 zip_path, output_path = task 

52 Path(output_path).write_bytes(convert_zip_to_nquads(zip_path)) 

53 return output_path 

54 

55 

56def open_output_file( 

57 output_dir: Path, prefix: str, file_index: int, compress: bool 

58) -> gzip.GzipFile | IO[bytes]: 

59 extension = "nq.gz" if compress else "nq" 

60 output_path = output_dir / f"{prefix}.{file_index:06d}.{extension}" 

61 if compress: 

62 return gzip.open(output_path, "wb") 

63 return output_path.open("wb") 

64 

65 

66def split_nquads_results(results: Iterable[bytes]) -> Iterator[list[bytes]]: 

67 for result in results: 

68 yield result.splitlines(keepends=True) 

69 

70 

71def read_nquads_file_groups(result_paths: Iterable[str]) -> Iterator[Iterator[bytes]]: 

72 for result_path in result_paths: 

73 path = Path(result_path) 

74 with path.open("rb") as file: 

75 yield iter(file.readline, b"") 

76 path.unlink() 

77 

78 

79def write_nquads_line_groups( 

80 line_groups: Iterable[Iterable[bytes]], 

81 output_dir: Path, 

82 prefix: str, 

83 lines_per_file: int, 

84 compress: bool, 

85 progress: Progress | None = None, 

86 task_id: TaskID | None = None, 

87) -> None: 

88 output_dir.mkdir(parents=True, exist_ok=True) 

89 file_index = 0 

90 line_count = 0 

91 output_file: gzip.GzipFile | IO[bytes] | None = None 

92 

93 try: 

94 for lines in line_groups: 

95 for line in lines: 

96 if output_file is None or line_count == lines_per_file: 

97 if output_file is not None: 

98 output_file.close() 

99 output_file = open_output_file( 

100 output_dir, prefix, file_index, compress 

101 ) 

102 file_index += 1 

103 line_count = 0 

104 output_file.write(line) 

105 line_count += 1 

106 if progress is not None and task_id is not None: 

107 progress.advance(task_id) 

108 finally: 

109 if output_file is not None: 

110 output_file.close() 

111 

112 

113def write_nquads_chunks( 

114 results: Iterable[bytes], 

115 output_dir: Path, 

116 prefix: str, 

117 lines_per_file: int, 

118 compress: bool, 

119 progress: Progress | None = None, 

120 task_id: TaskID | None = None, 

121) -> None: 

122 write_nquads_line_groups( 

123 split_nquads_results(results), 

124 output_dir, 

125 prefix, 

126 lines_per_file, 

127 compress, 

128 progress, 

129 task_id, 

130 ) 

131 

132 

133def write_nquads_stdout(results: Iterable[bytes]) -> None: 

134 stdout = sys.stdout.buffer 

135 for result in results: 

136 stdout.write(result) 

137 stdout.flush() 

138 

139 

140def create_progress() -> Progress: 

141 return Progress( 

142 SpinnerColumn(), 

143 TextColumn("[progress.description]{task.description}"), 

144 BarColumn(), 

145 MofNCompleteColumn(), 

146 TimeElapsedColumn(), 

147 TimeRemainingColumn(), 

148 console=Console(stderr=True), 

149 ) 

150 

151 

152def main() -> None: # pragma: no cover 

153 parser = argparse.ArgumentParser( 

154 description="Streams or writes chunked N-Quads from JSON-LD ZIP archives.", 

155 formatter_class=RichHelpFormatter, 

156 ) 

157 parser.add_argument( 

158 "rdf_dir", type=str, help="Root directory containing RDF ZIP archives" 

159 ) 

160 parser.add_argument( 

161 "-m", 

162 "--mode", 

163 type=str, 

164 choices=["all", "data", "prov"], 

165 default="all", 

166 help="Mode: 'all' for all ZIP files (default), 'data' for entity data only, 'prov' for provenance only", 

167 ) 

168 parser.add_argument( 

169 "-w", 

170 "--workers", 

171 type=int, 

172 default=None, 

173 help="Number of worker processes (defaults to min(8, CPU count))", 

174 ) 

175 parser.add_argument( 

176 "-o", 

177 "--output-dir", 

178 type=str, 

179 default=None, 

180 help="Directory where chunked N-Quads files are written instead of stdout", 

181 ) 

182 parser.add_argument( 

183 "--lines-per-file", 

184 type=int, 

185 default=DEFAULT_LINES_PER_FILE, 

186 help="Number of N-Quads lines per output file when --output-dir is used", 

187 ) 

188 parser.add_argument( 

189 "--gzip", 

190 action="store_true", 

191 help="Compress chunked output files with gzip when --output-dir is used", 

192 ) 

193 parser.add_argument( 

194 "--prefix", 

195 type=str, 

196 default="output", 

197 help="Output file prefix when --output-dir is used", 

198 ) 

199 args = parser.parse_args() 

200 

201 rdf_path = Path(args.rdf_dir).resolve() 

202 output_dir = Path(args.output_dir).resolve() if args.output_dir else None 

203 num_workers = args.workers if args.workers else min(8, multiprocessing.cpu_count()) 

204 

205 zip_files = collect_zip_files( 

206 str(rdf_path), 

207 only_data=args.mode == "data", 

208 only_prov=args.mode == "prov", 

209 ) 

210 

211 ctx = multiprocessing.get_context("forkserver") 

212 with ctx.Pool(processes=num_workers) as pool: 

213 if output_dir: 

214 output_dir.mkdir(parents=True, exist_ok=True) 

215 with create_progress() as progress: 

216 task_id = progress.add_task( 

217 "Writing N-Quads files", total=len(zip_files) 

218 ) 

219 with TemporaryDirectory( 

220 prefix=f".{args.prefix}.", dir=output_dir 

221 ) as temp_dir: 

222 temp_path = Path(temp_dir) 

223 conversion_tasks = ( 

224 (zip_path, str(temp_path / f"{index:06d}.nq")) 

225 for index, zip_path in enumerate(zip_files) 

226 ) 

227 result_paths = pool.imap_unordered( 

228 convert_zip_to_nquads_file, conversion_tasks, chunksize=10 

229 ) 

230 write_nquads_line_groups( 

231 read_nquads_file_groups(result_paths), 

232 output_dir, 

233 args.prefix, 

234 args.lines_per_file, 

235 args.gzip, 

236 progress, 

237 task_id, 

238 ) 

239 else: 

240 results = pool.imap_unordered( 

241 convert_zip_to_nquads, zip_files, chunksize=10 

242 ) 

243 write_nquads_stdout(results) 

244 

245 

246if __name__ == "__main__": # pragma: no cover 

247 main()