#!/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()