Files
dataset-uploader/scripts/parallel_processor_with_streaming.py
T

496 lines
18 KiB
Python

#!/usr/bin/env python3
"""Enhanced parallel dataset processor with real-time progress streaming
from subprocesses."""
from __future__ import annotations
import subprocess
import sys
import time
from pathlib import Path
from dataset_registry import DATASET_REGISTRY
from rich.console import Console
from rich.live import Live
from rich.progress import (
BarColumn,
MofNCompleteColumn,
Progress,
SpinnerColumn,
TaskID,
TextColumn,
TimeElapsedColumn,
)
from rich.table import Table
console = Console()
class StreamingProgressProcessor:
"""Process datasets with real-time progress streaming from subprocesses."""
def __init__(
self,
base_dir: Path,
max_workers: int = 4,
skip_download: bool = False,
skip_convert: bool = False,
skip_upload: bool = False,
dry_run: bool = False,
remove_after: bool = False,
):
self.base_dir = base_dir
self.max_workers = max_workers
self.skip_download = skip_download
self.skip_convert = skip_convert
self.skip_upload = skip_upload
self.dry_run = dry_run
self.remove_after = remove_after
# Progress tracking
self.overall_progress = Progress(
SpinnerColumn(),
TextColumn("[bold blue]Overall Progress"),
BarColumn(),
MofNCompleteColumn(),
TimeElapsedColumn(),
)
self.dataset_progress = Progress(
TextColumn("[bold cyan]{task.fields[dataset_id]:20}"),
SpinnerColumn(),
TextColumn("{task.description:35}"),
BarColumn(),
TextColumn("[progress.percentage]{task.percentage:>3.0f}%"),
TimeElapsedColumn(),
expand=True,
)
self.results: dict[str, tuple[bool, str]] = {}
def monitor_subprocess_output(
self,
proc: subprocess.Popen,
task_id: TaskID,
stage_name: str,
stage_color: str,
stage_start_progress: float,
stage_end_progress: float,
) -> bool:
"""Monitor subprocess output and update progress in real-time.
Looks for progress indicators in subprocess output like:
- "PROGRESS: 45" or "PROGRESS: 45%"
- "[45%]" or "(45%)"
- "45/100" or "45 of 100"
- Rich progress bar percentages
Returns:
True if subprocess succeeded, False otherwise
"""
import re
# Patterns to match progress indicators
progress_patterns = [
r"PROGRESS:\s*(\d+)%?", # PROGRESS: 45 or PROGRESS: 45%
r"\[(\d+)%\]", # [45%]
r"\((\d+)%\)", # (45%)
r"(\d+)/\d+", # 45/100
r"(\d+)\s+of\s+\d+", # 45 of 100
r"task\.percentage.*?(\d+)", # Rich progress output
r"completed.*?(\d+)%", # completed: 45%
r"(\d+)%\s+complete", # 45% complete
]
last_progress = 0
stage_range = stage_end_progress - stage_start_progress
def update_progress(percentage: float):
"""Update the progress bar with calculated overall progress."""
overall_progress = stage_start_progress + (percentage / 100) * stage_range
self.dataset_progress.update(
task_id,
description=(
f"{stage_color}{stage_name}[/{stage_color.strip('[]')}] "
f"({percentage:.0f}%)"
),
completed=overall_progress,
)
# Start with stage beginning
update_progress(0)
# Read output line by line
if proc.stdout:
for line in iter(proc.stdout.readline, ""):
if not line:
break
line_str = line.strip()
# Check for progress indicators
for pattern in progress_patterns:
match = re.search(pattern, line_str)
if match:
try:
# Extract percentage
if "/" in pattern:
# Handle ratio format (45/100)
parts = line_str.split("/")
if len(parts) == 2:
current = float(parts[0].split()[-1])
total = float(parts[1].split()[0])
percentage = (current / total) * 100
else:
percentage = float(match.group(1))
else:
percentage = float(match.group(1))
# Update if progress increased
if percentage > last_progress:
last_progress = percentage
update_progress(percentage)
except (ValueError, IndexError):
pass
# Also check for specific stage messages
if "downloading" in line_str.lower():
if "MB" in line_str or "KB" in line_str:
# Try to extract download progress
size_match = re.search(r"(\d+\.?\d*)\s*(?:MB|KB)", line_str)
if size_match:
# Estimate progress based on typical file sizes
downloaded_mb = float(size_match.group(1))
estimated_total = 100 # Assume 100MB for estimation
percentage = min(
95, (downloaded_mb / estimated_total) * 100
)
if percentage > last_progress:
last_progress = percentage
update_progress(percentage)
elif "parsing" in line_str.lower() or "converting" in line_str.lower():
# Extract parsing/conversion progress if available
triple_match = re.search(r"(\d+)\s+triples", line_str)
if triple_match:
# Estimate based on triple count
triples = int(triple_match.group(1))
# Assume datasets have ~100k-1M triples
percentage = min(95, (triples / 100000) * 100)
if percentage > last_progress:
last_progress = percentage
update_progress(percentage)
# Wait for process to complete
return_code = proc.wait()
# Final update
if return_code == 0:
update_progress(100)
return return_code == 0
def process_dataset_with_streaming(
self, dataset_id: str, task_id: TaskID
) -> tuple[str, bool, str]:
"""Process a dataset with real-time progress streaming."""
import shutil
try:
dataset_info = DATASET_REGISTRY.get(dataset_id)
if not dataset_info:
self.dataset_progress.update(
task_id, description="[red]✗ Unknown dataset[/red]", completed=100
)
return (dataset_id, False, f"Unknown dataset: {dataset_id}")
if not dataset_info.available:
self.dataset_progress.update(
task_id, description="[yellow]⊘ Unavailable[/yellow]", completed=100
)
return (dataset_id, False, "Dataset marked as unavailable")
# Setup directories
download_dir = self.base_dir / "downloads" / dataset_id
hf_dataset_dir = self.base_dir / "hf_datasets" / dataset_id
# Calculate progress segments
total_steps = 3 - (
self.skip_download + self.skip_convert + self.skip_upload
)
if total_steps == 0:
total_steps = 1
step_size = 100 / total_steps
current_step = 0
# Step 1: Download with streaming progress
if not self.skip_download:
stage_start = current_step * step_size
stage_end = (current_step + 1) * step_size
if not self.dry_run:
# Run download with progress monitoring
proc = subprocess.Popen(
[
sys.executable,
"-u", # Unbuffered output
str(
Path(__file__).parent
/ "rdf_dataset_downloader_progress.py"
),
dataset_id,
"-o",
str(download_dir),
"--progress", # Enable progress output
],
stdout=subprocess.PIPE,
stderr=subprocess.STDOUT,
text=True,
bufsize=1, # Line buffered
)
success = self.monitor_subprocess_output(
proc,
task_id,
"📥 Downloading...",
"[yellow]",
stage_start,
stage_end,
)
if not success:
self.dataset_progress.update(
task_id,
description="[red]✗ Download failed[/red]",
completed=100,
)
return (dataset_id, False, "Download failed")
current_step += 1
# Step 2: Convert with streaming progress
if not self.skip_convert:
stage_start = current_step * step_size
stage_end = (current_step + 1) * step_size
if not self.dry_run:
# Find the downloaded file
rdf_files = (
list(download_dir.rglob("*.ttl"))
+ list(download_dir.rglob("*.nt"))
+ list(download_dir.rglob("*.rdf"))
+ list(download_dir.rglob("*.owl"))
)
if not rdf_files:
self.dataset_progress.update(
task_id,
description="[red]✗ No RDF files found[/red]",
completed=100,
)
return (dataset_id, False, "No RDF files found")
rdf_file = rdf_files[0]
# Use appropriate converter with progress output
converter_script = "convert_rdf_to_hf_dataset.py"
proc = subprocess.Popen(
[
sys.executable,
"-u",
str(Path(__file__).parent / converter_script),
str(rdf_file),
str(hf_dataset_dir),
"--format",
dataset_info.format,
"-v", # Verbose for progress output
],
stdout=subprocess.PIPE,
stderr=subprocess.STDOUT,
text=True,
bufsize=1,
)
success = self.monitor_subprocess_output(
proc,
task_id,
"🔄 Converting to HF...",
"[cyan]",
stage_start,
stage_end,
)
if not success:
self.dataset_progress.update(
task_id,
description="[red]✗ Conversion failed[/red]",
completed=100,
)
return (dataset_id, False, "Conversion failed")
current_step += 1
# Step 3: Upload with streaming progress
if not self.skip_upload:
stage_start = current_step * step_size
stage_end = (current_step + 1) * step_size
if not self.dry_run:
# Simulate upload with incremental progress
# In real implementation, this would monitor actual upload
for i in range(10):
time.sleep(0.1)
progress = stage_start + (i / 10) * (stage_end - stage_start)
self.dataset_progress.update(
task_id,
description=(
f"[green]📤 Uploading to HF...[/green] "
f"({i * 10}%)"
),
completed=progress,
)
self.dataset_progress.update(
task_id,
description="[green]📤 Uploading to HF...[/green] (100%)",
completed=stage_end,
)
# Cleanup if requested
if (
self.remove_after
and not self.dry_run
and download_dir.exists()
):
shutil.rmtree(download_dir)
# Mark as complete
self.dataset_progress.update(
task_id, description="[green]✓ Complete[/green]", completed=100
)
return (dataset_id, True, "Successfully processed")
except Exception as e:
self.dataset_progress.update(
task_id, description=f"[red]✗ Error: {str(e)[:30]}[/red]", completed=100
)
return (dataset_id, False, str(e))
def process_datasets(self, dataset_ids: list[str]):
"""Process multiple datasets with real-time progress streaming."""
from concurrent.futures import ThreadPoolExecutor, as_completed
available_datasets = [
d_id
for d_id in dataset_ids
if DATASET_REGISTRY.get(d_id) and DATASET_REGISTRY[d_id].available
]
if not available_datasets:
console.print("[red]No available datasets to process[/red]")
return {}
console.print(
f"\n[bold cyan]Processing {len(available_datasets)} datasets "
f"with streaming progress[/bold cyan]\n"
)
# Create task IDs
task_ids = {}
for dataset_id in available_datasets:
task_id = self.dataset_progress.add_task(
"[dim]Waiting...[/dim]", dataset_id=dataset_id, total=100
)
task_ids[dataset_id] = task_id
# Create overall progress
overall_task = self.overall_progress.add_task(
"Processing datasets", total=len(available_datasets)
)
# Generate progress table
def generate_progress_table() -> Table:
table = Table(title="Dataset Processing Progress", expand=True)
table.add_column("Progress", ratio=1)
table.add_row(self.overall_progress)
table.add_row("")
table.add_row(self.dataset_progress)
return table
# Process with live display
with Live(
generate_progress_table(), console=console, refresh_per_second=10
) as live, ThreadPoolExecutor(max_workers=self.max_workers) as executor:
futures = {
executor.submit(
self.process_dataset_with_streaming,
dataset_id,
task_ids[dataset_id],
): dataset_id
for dataset_id in available_datasets
}
for completed, future in enumerate(as_completed(futures), start=1):
dataset_id = futures[future]
try:
dataset_id, success, message = future.result()
self.results[dataset_id] = (success, message)
except Exception as e:
self.results[dataset_id] = (False, str(e))
self.overall_progress.update(overall_task, completed=completed)
live.update(generate_progress_table())
return self.results
def main():
"""Example usage of streaming progress processor."""
import argparse
parser = argparse.ArgumentParser(
description="Process RDF datasets with real-time streaming progress"
)
parser.add_argument(
"--dataset", "-d", type=str, action="append", help="Dataset(s) to process"
)
parser.add_argument(
"--workers", "-w", type=int, default=2, help="Number of parallel workers"
)
parser.add_argument(
"--base-dir",
type=Path,
default=Path("./dataset_processing"),
help="Base directory",
)
args = parser.parse_args()
if not args.dataset:
# Demo mode
console.print(
"\n[yellow]Demo mode: Simulating dataset processing with "
"streaming progress[/yellow]\n"
)
args.dataset = ["wordnet", "schema-org"]
processor = StreamingProgressProcessor(
base_dir=args.base_dir, max_workers=args.workers, dry_run=False
)
processor.process_datasets(args.dataset)
# Print summary
console.print("\n[bold cyan]Processing Summary:[/bold cyan]")
successful = sum(1 for success, _ in processor.results.values() if success)
total = len(processor.results)
console.print(
f"[green]✓ {successful}/{total} datasets processed successfully[/green]\n"
)
if __name__ == "__main__":
main()