102 lines
3.4 KiB
Python
Executable File
102 lines
3.4 KiB
Python
Executable File
#!/usr/bin/env python3
|
|
"""
|
|
Wrapper for convert_rdf_to_hf_dataset_streaming_parallel.py that adds progress output
|
|
for subprocess monitoring in parallel processing.
|
|
"""
|
|
|
|
import subprocess
|
|
import sys
|
|
from pathlib import Path
|
|
from typing import Any
|
|
|
|
from scripts.common import print_progress
|
|
|
|
|
|
def monitor_and_output_progress(proc: subprocess.Popen[Any]) -> None:
|
|
"""Monitor subprocess output and emit simple progress markers."""
|
|
progress_patterns = [
|
|
# Streaming converter specific patterns
|
|
(
|
|
r"Chunk (\d+)/(\d+)",
|
|
lambda m: int(float(m.group(1)) / float(m.group(2))) * 70
|
|
if m.group(2) != "0"
|
|
else 0,
|
|
), # Chunk progress
|
|
(
|
|
r"Processing chunk (\d+)",
|
|
lambda m: min(5 + int(m.group(1)) * 5, 70),
|
|
), # Processing chunks
|
|
(r"Worker \d+: bytes", lambda m: 10), # Worker assignment
|
|
(r"✓ File divided", lambda m: 15), # File divided
|
|
(r"Parsing chunks", lambda m: 20), # Starting parsing
|
|
(
|
|
r"Processed (\d+) triples",
|
|
lambda m: min(70, 20 + int(m.group(1) / 1000)),
|
|
), # Triple count
|
|
# Phase detection with incremental progress
|
|
(r"Phase 1:", lambda m: 5), # Starting phase 1
|
|
(r"Dividing file", lambda m: 10), # Dividing file
|
|
(r"Phase 2:", lambda m: 20), # Starting phase 2
|
|
(r"Phase 3:", lambda m: 75), # Starting phase 3
|
|
(r"Creating dataset", lambda m: 80), # Creating dataset
|
|
(r"Writing to Parquet", lambda m: 85), # Writing parquet
|
|
(r"Saving dataset", lambda m: 90), # Saving dataset
|
|
(r"Building dataset", lambda m: 92), # Building dataset
|
|
(r"Dataset saved", lambda m: 95), # Dataset saved
|
|
(r"✓.*Successfully|Complete|Done", lambda m: 100), # Complete markers
|
|
# Memory usage patterns (informational)
|
|
(r"Memory usage:", lambda m: None), # Just pass through
|
|
# Generic percentage patterns
|
|
(r"(\d+)%", lambda m: int(m.group(1))), # Direct percentage
|
|
(r"\[(\d+)%\]", lambda m: int(m.group(1))), # Bracketed percentage
|
|
(r"Progress: (\d+\.?\d*)%", lambda m: int(m.group(1))), # Progress: X%
|
|
]
|
|
|
|
last_progress = 0
|
|
last_emit_time = 0
|
|
min_emit_interval = 0.5 # Minimum time between progress emissions
|
|
|
|
assert proc.stdout is not None
|
|
for line in iter(proc.stdout.readline, ""):
|
|
if not line:
|
|
break
|
|
|
|
# Pass through the original output
|
|
print(line, end="", flush=True)
|
|
|
|
# Try to extract and emit progress
|
|
|
|
last_progress = print_progress(
|
|
last_emit_time, last_progress, line, min_emit_interval, progress_patterns
|
|
)
|
|
|
|
# Ensure we hit 100% on successful completion
|
|
if proc.wait() == 0 and last_progress < 100:
|
|
print("\nPROGRESS: 100", flush=True)
|
|
|
|
|
|
def main() -> int:
|
|
"""Run the streaming converter with progress monitoring."""
|
|
# Pass all arguments to the actual converter
|
|
with subprocess.Popen(
|
|
[
|
|
sys.executable,
|
|
"-u",
|
|
str(
|
|
Path(__file__).parent
|
|
/ "convert_rdf_to_hf_dataset_streaming_parallel.py"
|
|
),
|
|
*sys.argv[1:],
|
|
],
|
|
stdout=subprocess.PIPE,
|
|
stderr=subprocess.STDOUT,
|
|
text=True,
|
|
bufsize=1,
|
|
) as proc:
|
|
monitor_and_output_progress(proc)
|
|
return proc.returncode
|
|
|
|
|
|
if __name__ == "__main__":
|
|
sys.exit(main())
|