diff --git a/owilix/cmd/subcmds/query_warc/parquet_logger.py b/owilix/cmd/subcmds/query_warc/parquet_logger.py index ebe03f9..670fc25 100644 --- a/owilix/cmd/subcmds/query_warc/parquet_logger.py +++ b/owilix/cmd/subcmds/query_warc/parquet_logger.py @@ -759,6 +759,85 @@ class ParquetJobLogger(object): "distributed_processes": len(self._cached_stats.file_mod_times) } + def store_completion_report(self, markdown_report: str) -> None: + """Store the completion report markdown for later retrieval.""" + try: + # Store the markdown report as a file alongside the parquet logs + report_filename = f"completion_report_{datetime.now().strftime('%Y%m%d_%H%M%S')}.md" + report_path = f"{self.destination_path}/{report_filename}" + + # Write the markdown report + with self.fs.open(report_path, "w", encoding="utf-8") as f: + f.write(markdown_report) + + # Store reference to the latest report + self._latest_report_path = report_path + self._latest_report_content = markdown_report + + if hasattr(self, 'verbose') and self.verbose: + print(f"πŸ“Š Completion report stored: {report_path}") + + except Exception as e: + print(f"⚠️ Could not store completion report: {e}") + # Store in memory as fallback + self._latest_report_content = markdown_report + + def get_completion_report_markdown(self) -> Optional[str]: + """Retrieve the stored completion report markdown.""" + try: + # Try to get from memory first + if hasattr(self, '_latest_report_content') and self._latest_report_content: + return self._latest_report_content + + # Try to read from the latest report file + if hasattr(self, '_latest_report_path') and self._latest_report_path: + try: + with self.fs.open(self._latest_report_path, "r", encoding="utf-8") as f: + return f.read() + except Exception: + pass + + # Look for the most recent report file + try: + # List all completion report files + files = self.fs.glob(f"{self.destination_path}/completion_report_*.md") + if files: + # Sort by filename (which includes timestamp) and get the latest + latest_file = sorted(files)[-1] + with self.fs.open(latest_file, "r", encoding="utf-8") as f: + report_content = f.read() + # Cache for future use + self._latest_report_content = report_content + self._latest_report_path = latest_file + return report_content + except Exception: + pass + + return None + + except Exception as e: + print(f"⚠️ Could not retrieve completion report: {e}") + return None + + def list_completion_reports(self) -> List[str]: + """List all stored completion reports.""" + try: + files = self.fs.glob(f"{self.destination_path}/completion_report_*.md") + return sorted(files) + except Exception as e: + print(f"⚠️ Could not list completion reports: {e}") + return [] + + def get_completion_report_by_timestamp(self, timestamp: str) -> Optional[str]: + """Retrieve a completion report by timestamp (YYYYMMDD_HHMMSS format).""" + try: + report_path = f"{self.destination_path}/completion_report_{timestamp}.md" + with self.fs.open(report_path, "r", encoding="utf-8") as f: + return f.read() + except Exception as e: + print(f"⚠️ Could not retrieve completion report for timestamp {timestamp}: {e}") + return None + def cleanup_stale_files(self, max_age_hours: int = 168): # 1 week default """ Clean up stale process files from crashed or terminated processes. @@ -798,73 +877,25 @@ class ParquetJobLogger(object): self.flush() - -class TransactionLogger: - """Transaction logger for file processing results in JSON Lines format.""" - - def __init__(self, log_path: str, enabled: bool = True): - self.log_path = log_path - self.enabled = enabled - self._lock = threading.Lock() - - if self.enabled and log_path: - os.makedirs(os.path.dirname(log_path), exist_ok=True) - - def log_file_result(self, warc_file: str, source_key: str, total_tasks: int, - successful: int, failed: int, processing_time: float, - error_details: List[str] = None): - """Log comprehensive results for a file processing job.""" - if not self.enabled or not self.log_path: - return - - with self._lock: - entry = { - "timestamp": datetime.now().isoformat(), - "warc_file": warc_file, - "source_key": source_key, - "total_tasks": total_tasks, - "successful_records": successful, - "failed_records": failed, - "processing_time_seconds": processing_time, - "success_rate": (successful / max(total_tasks, 1)) * 100, - "error_details": error_details or [] - } - self._write_entry(entry) - - def _write_entry(self, entry: dict): - try: - with open(self.log_path, "a", encoding="utf-8") as f: - f.write(json.dumps(entry) + "\n") - except Exception as e: - print(f"Failed to write to transaction log: {e}") - - -def analyze_warc_log(self, dummy_local, dummy_remote, log_directory: str = None, +def analyze_warc_log(self, dummy_local, dummy_remote, log: str = None, verbose: bool = False, show_failed_jobs: bool = False, show_duplicates: bool = False, cleanup_old_entries: bool = False, days_to_keep: int = 30, consolidate_logs: bool = False): """ - Analyze distributed parquet job logs from WARC processing with enhanced multiprocess support. - + Analyze distributed parquet job logs from WARC cache fetchs. This enhanced analysis function handles the new distributed logging architecture where each process maintains its own log file in a shared directory. It provides comprehensive statistics across all processes with performance insights, load balancing analysis, and distributed system health monitoring. CLI Usage: - owilix query analyze_distributed_warc_log log_directory=/shared/logs + owilix query analyze_distributed_warc_log log=/shared/logs + - Key Features: - - **Distributed Log Analysis**: Aggregates statistics across multiple process log files - - **Process Performance Monitoring**: Per-process throughput and efficiency analysis - - **Load Balancing Assessment**: Workload distribution analysis across processes - - **System Health Metrics**: Detect stale processes, crashed workers, and bottlenecks - - **Temporal Analysis**: Processing patterns over time with distributed context - - **Fault Tolerance**: Handles corrupted files and missing data gracefully Args: self: Command processor instance (for accessing console and utilities) - log_directory (str): Path to directory containing distributed parquet log files + log (str): Path to directory containing distributed parquet log files verbose (bool): Enable detailed per-process and per-file information show_failed_jobs (bool): Display detailed failure analysis across processes show_duplicates (bool): Analyze duplicate processing across processes @@ -885,25 +916,25 @@ def analyze_warc_log(self, dummy_local, dummy_remote, log_directory: str = None, try: # Validate input directory - if not log_directory or not os.path.exists(log_directory): + if not log or not os.path.exists(log): return CommandResult( success=False, object={}, - msg=f"Log directory not found: {log_directory}" + msg=f"Log directory not found: {log}" ) - self.console.print(f"[cyan]πŸ“Š Analyzing Distributed WARC Job Logs: {log_directory}[/cyan]") + self.console.print(f"[cyan]πŸ“Š Analyzing Distributed WARC Job Logs: {log}[/cyan]") # Discover all log files log_files = [] - for file_path in os.listdir(log_directory): + for file_path in os.listdir(log): if file_path.endswith('.parquet') and not file_path.endswith('_temp.parquet'): - log_files.append(os.path.join(log_directory, file_path)) + log_files.append(os.path.join(log, file_path)) if not log_files: return CommandResult( success=True, - object={"log_directory": log_directory, "total_files": 0}, + object={"log": log, "total_files": 0}, msg="No log files found in directory" ) @@ -947,7 +978,7 @@ def analyze_warc_log(self, dummy_local, dummy_remote, log_directory: str = None, # Combine all dataframes combined_df = pd.concat(all_dataframes, ignore_index=True) - combined_df['timestamp_dt'] = pd.to_datetime(combined_df['timestamp']) + combined_df['timestamp_dt'] = pd.to_datetime(combined_df['timestamp'], format='mixed') total_entries = len(combined_df) unique_jobs = combined_df['job_id'].nunique() @@ -979,7 +1010,7 @@ def analyze_warc_log(self, dummy_local, dummy_remote, log_directory: str = None, 'jobs_completed': len(completed_df), 'total_tasks': completed_df['task_count'].sum(), 'avg_processing_time': completed_df['processing_time_seconds'].mean(), - 'total_bytes_written': completed_df['total_bytes_written'].sum(), + 'file_size': completed_df['file_size'].sum(), 'avg_seek_time_ms': completed_df['avg_seek_time_ms'].mean(), 'avg_read_time_ms': completed_df['avg_read_time_ms'].mean(), 'avg_write_time_ms': completed_df['avg_write_time_ms'].mean(), @@ -991,7 +1022,7 @@ def analyze_warc_log(self, dummy_local, dummy_remote, log_directory: str = None, # Standard performance metrics (same as original) completed_jobs = latest_status_df[latest_status_df['status'] == 'completed'] total_tasks_completed = completed_jobs['task_count'].sum() - total_bytes_written = completed_jobs['total_bytes_written'].sum() + file_size = completed_jobs['file_size'].sum() # Enhanced timing analysis avg_seek_time = completed_jobs['avg_seek_time_ms'].mean() if len(completed_jobs) > 0 else 0 @@ -1005,8 +1036,8 @@ def analyze_warc_log(self, dummy_local, dummy_remote, log_directory: str = None, success_rate = (completed_count / max(unique_jobs, 1)) * 100 # Display enhanced distributed analysis - self.console.print(f"\n[bold green]πŸ“‹ Distributed WARC Job Log Analysis Report[/bold green]") - self.console.print(f"[cyan]πŸ“‚ Directory: {log_directory}[/cyan]") + self.console.print(f"\n[bold green]πŸ“‹ WARC Cache Log Analysis Report[/bold green]") + self.console.print(f"[cyan]πŸ“‚ Directory: {log}[/cyan]") # Distributed system overview self.console.print(f"\n[bold]πŸ—οΈ Distributed System Overview:[/bold]") @@ -1037,14 +1068,14 @@ def analyze_warc_log(self, dummy_local, dummy_remote, log_directory: str = None, if completed_count > 0: self.console.print(f"\n[bold]⚑ Aggregated Performance Metrics:[/bold]") self.console.print(f" β€’ Total tasks completed: {total_tasks_completed:,}") - self.console.print(f" β€’ Total bytes written: {total_bytes_written:,} ({total_bytes_written / (1024**3):.2f} GB)") + self.console.print(f" β€’ Total bytes written: {file_size:,} ({file_size / (1024**3):.2f} GB)") self.console.print(f" β€’ Average timing per task:") self.console.print(f" - Seek: {avg_seek_time:.2f}ms") self.console.print(f" - Read: {avg_read_time:.2f}ms") self.console.print(f" - Write: {avg_write_time:.2f}ms") if processing_duration > 0: - throughput_gb_hour = (total_bytes_written / (1024**3)) / processing_duration + throughput_gb_hour = (file_size / (1024**3)) / processing_duration self.console.print(f" β€’ System throughput: {throughput_gb_hour:.2f} GB/hour") # Process-level performance analysis @@ -1118,33 +1149,11 @@ def analyze_warc_log(self, dummy_local, dummy_remote, log_directory: str = None, else: self.console.print(f" β€’ [red]⚠️ High performance variation ({cv:.1f}% CV)[/red]") - # Enhanced recommendations - self.console.print(f"\n[bold]πŸ’‘ Distributed System Recommendations:[/bold]") - - if success_rate > 95: - self.console.print(f" β€’ [green]βœ… Excellent success rate ({success_rate:.1f}%)[/green]") - elif success_rate > 85: - self.console.print(f" β€’ [yellow]βš–οΈ Good success rate ({success_rate:.1f}%) - monitor failures[/yellow]") - else: - self.console.print(f" β€’ [red]⚠️ Low success rate ({success_rate:.1f}%) - investigate systematic issues[/red]") - - if len(stale_processes) > 0: - self.console.print(f" β€’ [yellow]🧹 Consider cleanup of {len(stale_processes)} stale process files[/yellow]") - - if len(throughputs) > 1 and (throughput_stddev / max(avg_throughput, 1)) * 100 > 30: - self.console.print(f" β€’ [red]βš–οΈ Investigate process performance disparities[/red]") - - self.console.print(f" β€’ [cyan]πŸ“Š {len(all_dataframes)} distributed log files processed successfully[/cyan]") - self.console.print(f" β€’ [cyan]πŸ”„ Consider periodic log consolidation for long-term storage[/cyan]") - # Analysis timing - analysis_time = time.time() - start_time - self.console.print(f"\n[dim]Distributed analysis completed in {analysis_time:.2f} seconds[/dim]") # Prepare comprehensive result object result_object = { - "analysis_time": analysis_time, - "log_directory": log_directory, + "log": log, "distributed_system": { "total_log_files": len(log_files), "loaded_files": len(all_dataframes), @@ -1171,11 +1180,11 @@ def analyze_warc_log(self, dummy_local, dummy_remote, log_directory: str = None, }, "performance_metrics": { "total_tasks_completed": int(total_tasks_completed), - "total_bytes_written": int(total_bytes_written), + "file_size": int(file_size), "avg_seek_time_ms": avg_seek_time, "avg_read_time_ms": avg_read_time, "avg_write_time_ms": avg_write_time, - "throughput_gb_per_hour": (total_bytes_written / (1024**3)) / max(processing_duration, 0.001) + "throughput_gb_per_hour": (file_size / (1024**3)) / max(processing_duration, 0.001) }, "process_performance": process_performance, "load_balancing": { @@ -1197,11 +1206,13 @@ def analyze_warc_log(self, dummy_local, dummy_remote, log_directory: str = None, msg=f"Successfully analyzed {len(all_dataframes)} distributed log files with {total_entries:,} entries. " f"Success rate: {success_rate:.1f}%, " f"Active processes: {active_processes}, " - f"System throughput: {(total_bytes_written / (1024**3)) / max(processing_duration, 0.001):.2f} GB/hour, " + f"System throughput: {(file_size / (1024**3)) / max(processing_duration, 0.001):.2f} GB/hour, " f"Load balance ratio: {min(job_counts) / max(max(job_counts), 1):.2f}" if job_counts else "0.0" ) except Exception as e: + import traceback + traceback.print_exc() self.console.print(f"[red]❌ Distributed job log analysis failed: {e}[/red]") return CommandResult( success=False, diff --git a/owilix/cmd/subcmds/query_warc/query_warc.py b/owilix/cmd/subcmds/query_warc/query_warc.py index 7a25987..11781c4 100644 --- a/owilix/cmd/subcmds/query_warc/query_warc.py +++ b/owilix/cmd/subcmds/query_warc/query_warc.py @@ -160,7 +160,7 @@ from rich.text import Text from owilix.cmd.base import CommandResult from owilix.core.duckdb import OWIlixSQLQuery, OWIDuckDBSelectExecutor -from .parquet_logger import ParquetJobLogger, get_fs, TransactionLogger +from .parquet_logger import ParquetJobLogger, get_fs try: import zmq @@ -200,42 +200,32 @@ class DatacenterStats: """Statistics tracking for a single datacenter queue with corrected bandwidth analysis.""" datacenter: str start_time: float = field(default_factory=time.time) - # Queue status jobs_queued: int = 0 jobs_completed: int = 0 jobs_failed: int = 0 jobs_started: int = 0 - # Processing metrics total_processing_time: float = 0.0 total_tasks_processed: int = 0 total_successful_tasks: int = 0 total_failed_tasks: int = 0 - # Enhanced timing metrics with detailed breakdown - total_seek_time: float = 0.0 - total_read_time: float = 0.0 - total_write_time: float = 0.0 + avg_seek_time: float = 0.0 + avg_read_time: float = 0.0 + avg_write_time: float = 0.0 total_file_open_time: float = 0.0 + avg_bytes_per_write: float = 0.0 # CORRECTED: Single bytes tracking total_bytes: int = 0 # Total bytes processed - - # Individual timing measurements for standard deviation calculation (keep last 100) - seek_times: List[float] = field(default_factory=lambda: deque(maxlen=100)) - read_times: List[float] = field(default_factory=lambda: deque(maxlen=100)) - write_times: List[float] = field(default_factory=lambda: deque(maxlen=100)) - file_open_times: List[float] = field(default_factory=lambda: deque(maxlen=100)) - - # CORRECTED: Bandwidth measurements (keep last 100 for rolling average) - read_bandwidth_samples: List[float] = field(default_factory=lambda: deque(maxlen=100)) # MiB/s - write_bandwidth_samples: List[float] = field(default_factory=lambda: deque(maxlen=100)) # MiB/s - combined_bandwidth_samples: List[float] = field(default_factory=lambda: deque(maxlen=100)) # MiB/s - # Performance tracking last_completed_time: float = 0.0 - completion_times: List[float] = field(default_factory=list) + + def __post_init__(self): + # Ξ± = 2 / (N + 1) for an EMA spanning ~N samples + N = 500.0 + self._ema_alpha = 2.0 / (N + 1.0) def add_job_queued(self): """Record a job being added to the queue.""" @@ -256,53 +246,20 @@ class DatacenterStats: self.total_tasks_processed += successful_tasks + failed_tasks self.total_successful_tasks += successful_tasks self.total_failed_tasks += failed_tasks + Ξ± = self._ema_alpha + # Update EMAs: new_avg = Ξ± * new_sample + (1 - Ξ±) * old_avg + self.avg_seek_time = Ξ± * seek_time + (1 - Ξ±) * self.avg_seek_time + self.avg_read_time = Ξ± * read_time + (1 - Ξ±) * self.avg_read_time + self.avg_write_time = Ξ± * write_time + (1 - Ξ±) * self.avg_write_time + self.avg_bytes_per_write = Ξ± * bytes_processed + (1 - Ξ±) * self.avg_bytes_per_write - # Add detailed timing information - self.total_seek_time += seek_time - self.total_read_time += read_time - self.total_write_time += write_time self.total_file_open_time += file_open_time # CORRECTED: Single bytes tracking self.total_bytes += bytes_processed - # Track individual measurements for standard deviation - if seek_time > 0: - self.seek_times.append(seek_time / max(successful_tasks + failed_tasks, 1)) - if read_time > 0: - self.read_times.append(read_time / max(successful_tasks + failed_tasks, 1)) - if write_time > 0: - self.write_times.append(write_time / max(successful_tasks, 1)) - if file_open_time > 0: - self.file_open_times.append(file_open_time) - - # CORRECTED: Track bandwidth samples using corrected calculations - if bytes_processed > 0: - bytes_mib = bytes_processed / (1024 * 1024) - - # Read bandwidth = total_bytes / read_time - if read_time > 0: - read_bandwidth = bytes_mib / read_time - self.read_bandwidth_samples.append(read_bandwidth) - - # Write bandwidth = total_bytes / write_time - if write_time > 0: - write_bandwidth = bytes_mib / write_time - self.write_bandwidth_samples.append(write_bandwidth) - - # Combined bandwidth = 2 * total_bytes / (read_time + write_time) - total_io_time = read_time + write_time - if total_io_time > 0: - combined_bandwidth = (2 * bytes_mib) / total_io_time - self.combined_bandwidth_samples.append(combined_bandwidth) - current_time = time.time() self.last_completed_time = current_time - self.completion_times.append(current_time) - - # Keep only recent completion times for ETA calculation (last 10 completions) - if len(self.completion_times) > 10: - self.completion_times = self.completion_times[-10:] def add_job_failed(self, processing_time: float, task_count: int): """Record a job failure.""" @@ -314,22 +271,12 @@ class DatacenterStats: current_time = time.time() self.last_completed_time = current_time - self.completion_times.append(current_time) - if len(self.completion_times) > 10: - self.completion_times = self.completion_times[-10:] def get_completion_rate(self) -> float: """Get jobs completed per second based on recent completions.""" - if len(self.completion_times) < 2: - return 0.0 - - # Calculate rate based on recent completions - time_span = self.completion_times[-1] - self.completion_times[0] - if time_span <= 0: - return 0.0 - - return (len(self.completion_times) - 1) / time_span + _elapsed_time = time.time() - self.last_completed_time + return self.jobs_completed / _elapsed_time def get_eta_seconds(self) -> Optional[float]: """Calculate ETA for remaining jobs in seconds.""" @@ -362,145 +309,88 @@ class DatacenterStats: def get_timing_stats(self) -> dict: """Get detailed timing statistics with averages and standard deviations including corrected bandwidth.""" - completed_tasks = max(self.total_successful_tasks, 1) completed_jobs = max(self.jobs_completed, 1) - - # Calculate averages - avg_seek_ms = (self.total_seek_time / completed_tasks) * 1000 - avg_read_ms = (self.total_read_time / completed_tasks) * 1000 - avg_write_ms = (self.total_write_time / max(self.total_successful_tasks, 1)) * 1000 avg_file_open_ms = (self.total_file_open_time / completed_jobs) * 1000 - # Calculate standard deviations - seek_stddev = statistics.stdev([t * 1000 for t in self.seek_times]) if len(self.seek_times) > 1 else 0.0 - read_stddev = statistics.stdev([t * 1000 for t in self.read_times]) if len(self.read_times) > 1 else 0.0 - write_stddev = statistics.stdev([t * 1000 for t in self.write_times]) if len(self.write_times) > 1 else 0.0 - file_open_stddev = statistics.stdev([t * 1000 for t in self.file_open_times]) if len(self.file_open_times) > 1 else 0.0 - - # CORRECTED: Calculate bandwidth statistics using corrected formulas - avg_read_bandwidth = statistics.mean(self.read_bandwidth_samples) if self.read_bandwidth_samples else 0.0 - avg_write_bandwidth = statistics.mean(self.write_bandwidth_samples) if self.write_bandwidth_samples else 0.0 - avg_combined_bandwidth = statistics.mean(self.combined_bandwidth_samples) if self.combined_bandwidth_samples else 0.0 - - read_bandwidth_stddev = statistics.stdev(self.read_bandwidth_samples) if len(self.read_bandwidth_samples) > 1 else 0.0 - write_bandwidth_stddev = statistics.stdev(self.write_bandwidth_samples) if len(self.write_bandwidth_samples) > 1 else 0.0 - combined_bandwidth_stddev = statistics.stdev(self.combined_bandwidth_samples) if len(self.combined_bandwidth_samples) > 1 else 0.0 - # Overall bandwidth (based on total time and bytes) total_processing_time_sec = max(self.total_processing_time, 0.001) overall_bandwidth = (self.total_bytes / (1024 * 1024)) / total_processing_time_sec + _tb_mb =self.total_bytes / (1024 * 1024) + _elapsed_s = time.time() - self.start_time return { - "avg_seek_time_ms": avg_seek_ms, - "avg_read_time_ms": avg_read_ms, - "avg_write_time_ms": avg_write_ms, + "elapsed_s": _elapsed_s, + "avg_seek_time_ms": self.avg_seek_time*1e3, + "avg_read_time_ms": self.avg_read_time*1e3, + "avg_write_time_ms": self.avg_write_time*1e3, "avg_file_open_time_ms": avg_file_open_ms, - "seek_stddev_ms": seek_stddev, - "read_stddev_ms": read_stddev, - "write_stddev_ms": write_stddev, - "file_open_stddev_ms": file_open_stddev, # CORRECTED: Bandwidth statistics using proper calculations - "avg_read_bandwidth_mib_s": avg_read_bandwidth, - "avg_write_bandwidth_mib_s": avg_write_bandwidth, - "avg_combined_bandwidth_mib_s": avg_combined_bandwidth, - "read_bandwidth_stddev_mib_s": read_bandwidth_stddev, - "write_bandwidth_stddev_mib_s": write_bandwidth_stddev, - "combined_bandwidth_stddev_mib_s": combined_bandwidth_stddev, + "avg_read_bandwidth_mib_s": _tb_mb/(self.avg_read_time+1), + "avg_write_bandwidth_mib_s": _tb_mb/(self.avg_write_time+1), + "avg_bandwidth_mib_s": _tb_mb/_elapsed_s, "overall_bandwidth_mib_s": overall_bandwidth, - "total_bytes_mb": self.total_bytes / (1024 * 1024), - "measurements": { - "seek_samples": len(self.seek_times), - "read_samples": len(self.read_times), - "write_samples": len(self.write_times), - "file_open_samples": len(self.file_open_times), - "read_bandwidth_samples": len(self.read_bandwidth_samples), - "write_bandwidth_samples": len(self.write_bandwidth_samples), - "combined_bandwidth_samples": len(self.combined_bandwidth_samples) - } + "total_bytes_mb": _tb_mb, } - @dataclass class WARCDestinationStats: - """Performance statistics for a single WARC destination (per thread) with corrected timing and bandwidth analysis.""" + """Performance statistics for a single WARC destination (per thread) + with corrected timing and bandwidth analysis using exponential moving averages.""" thread_id: int creation_time: float = field(default_factory=time.time) - # File statistics records_written: int = 0 warc_files_created: int = 0 - total_bytes: int = 0 # CORRECTED: Single bytes tracking - - # Timing statistics - total_seek_time: float = 0.0 - total_read_time: float = 0.0 - total_write_time: float = 0.0 - total_file_open_time: float = 0.0 - - # Individual timing measurements for standard deviation calculation (keep last 100) - seek_times: deque = field(default_factory=lambda: deque(maxlen=100)) - read_times: deque = field(default_factory=lambda: deque(maxlen=100)) - write_times: deque = field(default_factory=lambda: deque(maxlen=100)) - file_open_times: deque = field(default_factory=lambda: deque(maxlen=100)) - - # CORRECTED: Bandwidth measurements (MiB/s) - keep last 100 for rolling average - read_bandwidth_samples: deque = field(default_factory=lambda: deque(maxlen=100)) - write_bandwidth_samples: deque = field(default_factory=lambda: deque(maxlen=100)) - combined_bandwidth_samples: deque = field(default_factory=lambda: deque(maxlen=100)) - + total_bytes: int = 0 + # Exponential moving averages (EMAs) of timings and sizes + avg_seek_time: float = 0.0 # seconds + avg_read_time: float = 0.0 # seconds + avg_write_time: float = 0.0 # seconds + avg_file_open_time: float = 0.0 # seconds + avg_bytes_per_write: float = 0.0 # bytes # Error tracking write_errors: int = 0 file_creation_errors: int = 0 - # Performance metrics last_write_time: float = 0.0 writes_per_second: float = 0.0 + # Internal: smoothing factor for EMAs (approx. window β‰ˆ 500) + _ema_alpha: float = field(init=False, repr=False) - def add_write_operation(self, seek_time: float, read_time: float, write_time: float, - bytes_processed: int): - """Record a successful write operation with corrected timing details and bandwidth tracking.""" + def __post_init__(self): + # Ξ± = 2 / (N + 1) for an EMA spanning ~N samples + N = 500.0 + self._ema_alpha = 2.0 / (N + 1.0) + + def add_write_operation(self, + seek_time: float, + read_time: float, + write_time: float, + bytes_processed: int): + """Record a successful write operation and update EMAs.""" + # Update counts and totals self.records_written += 1 - self.total_seek_time += seek_time - self.total_read_time += read_time - self.total_write_time += write_time - self.total_bytes += bytes_processed # CORRECTED: Single bytes tracking - self.last_write_time = time.time() - - # Store individual measurements for standard deviation calculation - self.seek_times.append(seek_time) - self.read_times.append(read_time) - self.write_times.append(write_time) - - # CORRECTED: Track bandwidth samples using corrected calculations - if bytes_processed > 0: - bytes_mib = bytes_processed / (1024 * 1024) - - # Read bandwidth = total_bytes / read_time - if read_time > 0: - read_bandwidth = bytes_mib / read_time - self.read_bandwidth_samples.append(read_bandwidth) - - # Write bandwidth = total_bytes / write_time - if write_time > 0: - write_bandwidth = bytes_mib / write_time - self.write_bandwidth_samples.append(write_bandwidth) - - # Combined bandwidth = 2 * total_bytes / (read_time + write_time) - total_io_time = read_time + write_time - if total_io_time > 0: - combined_bandwidth = (2 * bytes_mib) / total_io_time - self.combined_bandwidth_samples.append(combined_bandwidth) - - # Calculate writes per second (moving average) - elapsed = self.last_write_time - self.creation_time + self.total_bytes += bytes_processed + now = time.time() + self.last_write_time = now + + Ξ± = self._ema_alpha + # Update EMAs: new_avg = Ξ± * new_sample + (1 - Ξ±) * old_avg + self.avg_seek_time = Ξ± * seek_time + (1 - Ξ±) * self.avg_seek_time + self.avg_read_time = Ξ± * read_time + (1 - Ξ±) * self.avg_read_time + self.avg_write_time = Ξ± * write_time + (1 - Ξ±) * self.avg_write_time + self.avg_bytes_per_write = Ξ± * bytes_processed + (1 - Ξ±) * self.avg_bytes_per_write + + # Recalculate writes per second over lifetime + elapsed = now - self.creation_time if elapsed > 0: self.writes_per_second = self.records_written / elapsed def add_file_creation(self, creation_time: float): - """Record a new WARC file creation with timing.""" + """Record a new WARC file creation and update EMA of open time.""" self.warc_files_created += 1 - self.total_file_open_time += creation_time - self.file_open_times.append(creation_time) + Ξ± = self._ema_alpha + self.avg_file_open_time = Ξ± * creation_time + (1 - Ξ±) * self.avg_file_open_time def add_write_error(self): """Record a write error.""" @@ -511,90 +401,47 @@ class WARCDestinationStats: self.file_creation_errors += 1 def get_activity_level(self) -> str: - """Get worker activity level based on recent performance.""" - if self.writes_per_second > 10: + """Get worker activity level based on current writes/sec.""" + wps = self.writes_per_second + if wps > 10: return "Very Active" - elif self.writes_per_second > 5: + elif wps > 5: return "Active" - elif self.writes_per_second > 1: + elif wps > 1: return "Moderate" - elif self.writes_per_second > 0: + elif wps > 0: return "Low" else: return "Idle" def get_summary(self) -> dict: - """Get comprehensive statistics summary with corrected timing and bandwidth analysis.""" + """Return a comprehensive summary of all statistics.""" elapsed = time.time() - self.creation_time - avg_seek_time = (self.total_seek_time / max(self.records_written, 1)) * 1000 - avg_read_time = (self.total_read_time / max(self.records_written, 1)) * 1000 - avg_write_time = (self.total_write_time / max(self.records_written, 1)) * 1000 - avg_file_open_time = (self.total_file_open_time / max(self.warc_files_created, 1)) * 1000 - - # Calculate standard deviations from individual measurements - seek_stddev = statistics.stdev([t * 1000 for t in self.seek_times]) if len(self.seek_times) > 1 else 0.0 - read_stddev = statistics.stdev([t * 1000 for t in self.read_times]) if len(self.read_times) > 1 else 0.0 - write_stddev = statistics.stdev([t * 1000 for t in self.write_times]) if len(self.write_times) > 1 else 0.0 - file_open_stddev = statistics.stdev([t * 1000 for t in self.file_open_times]) if len(self.file_open_times) > 1 else 0.0 - - # CORRECTED: Calculate bandwidth statistics using corrected formulas - avg_read_bandwidth = statistics.mean(self.read_bandwidth_samples) if self.read_bandwidth_samples else 0.0 - avg_write_bandwidth = statistics.mean(self.write_bandwidth_samples) if self.write_bandwidth_samples else 0.0 - avg_combined_bandwidth = statistics.mean(self.combined_bandwidth_samples) if self.combined_bandwidth_samples else 0.0 - - read_bandwidth_stddev = statistics.stdev(self.read_bandwidth_samples) if len(self.read_bandwidth_samples) > 1 else 0.0 - write_bandwidth_stddev = statistics.stdev(self.write_bandwidth_samples) if len(self.write_bandwidth_samples) > 1 else 0.0 - combined_bandwidth_stddev = statistics.stdev(self.combined_bandwidth_samples) if len(self.combined_bandwidth_samples) > 1 else 0.0 - - # Overall bandwidth (total bytes / total time) - overall_bandwidth = (self.total_bytes / (1024 * 1024)) / max(elapsed, 0.001) + error_rate = ( + self.write_errors / max(self.records_written + self.write_errors, 1) + ) * 100.0 return { "thread_id": self.thread_id, - "elapsed_time": elapsed, + "elapsed_time_s": elapsed, "records_written": self.records_written, "warc_files_created": self.warc_files_created, "total_bytes": self.total_bytes, "writes_per_second": self.writes_per_second, "activity_level": self.get_activity_level(), - - # Enhanced timing with standard deviations - "avg_seek_time_ms": avg_seek_time, - "avg_read_time_ms": avg_read_time, - "avg_write_time_ms": avg_write_time, - "avg_file_open_time_ms": avg_file_open_time, - "seek_stddev_ms": seek_stddev, - "read_stddev_ms": read_stddev, - "write_stddev_ms": write_stddev, - "file_open_stddev_ms": file_open_stddev, - - # CORRECTED: Bandwidth statistics using proper calculations - "avg_read_bandwidth_mib_s": avg_read_bandwidth, - "avg_write_bandwidth_mib_s": avg_write_bandwidth, - "avg_combined_bandwidth_mib_s": avg_combined_bandwidth, - "read_bandwidth_stddev_mib_s": read_bandwidth_stddev, - "write_bandwidth_stddev_mib_s": write_bandwidth_stddev, - "combined_bandwidth_stddev_mib_s": combined_bandwidth_stddev, - "overall_bandwidth_mib_s": overall_bandwidth, - - "total_processing_time_ms": avg_seek_time + avg_read_time + avg_write_time, - "write_errors": self.write_errors, - "file_creation_errors": self.file_creation_errors, - "error_rate_percent": (self.write_errors / max(self.records_written + self.write_errors, 1)) * 100, - - # Measurement sample counts - "timing_samples": { - "seek_samples": len(self.seek_times), - "read_samples": len(self.read_times), - "write_samples": len(self.write_times), - "file_open_samples": len(self.file_open_times), - "read_bandwidth_samples": len(self.read_bandwidth_samples), - "write_bandwidth_samples": len(self.write_bandwidth_samples), - "combined_bandwidth_samples": len(self.combined_bandwidth_samples) - } + "avg_seek_time_ms": self.avg_seek_time * 1e3, + "avg_read_time_ms": self.avg_read_time * 1e3, + "avg_write_time_ms": self.avg_write_time * 1e3, + "avg_file_open_time_ms": self.avg_file_open_time * 1e3, + "avg_bytes_per_write": self.avg_bytes_per_write, + "write_errors": self.write_errors, + "file_creation_errors": self.file_creation_errors, + "error_rate_percent": error_rate, + "avg_read_bandwidth_mib_s": self.total_bytes / (1024 * 1024 * self.avg_read_time+1), + "avg_write_bandwidth_mib_s": self.total_bytes / (1024 * 1024 * self.avg_read_time+1), + "avg_bandwidth_mib_s": self.total_bytes / (1024 * 1024 * elapsed+1), } - class WARCDestination: """High-performance WARC destination writer for single-thread use with corrected bandwidth tracking.""" @@ -738,63 +585,52 @@ class ParallelWARCDestinationManager: def get_aggregated_stats(self) -> dict: """Get aggregated statistics across all destinations including corrected bandwidth.""" - all_stats = self.get_all_stats() - - if not all_stats: + stats_map = self.get_all_stats() + if not stats_map: return {} - # Aggregate metrics - total_records = sum(stats["records_written"] for stats in all_stats.values()) - total_files = sum(stats["warc_files_created"] for stats in all_stats.values()) - total_bytes = sum(stats["total_bytes"] for stats in all_stats.values()) - total_errors = sum(stats["write_errors"] for stats in all_stats.values()) - - # Calculate averages - active_threads = len(all_stats) - avg_writes_per_second = sum(stats["writes_per_second"] for stats in all_stats.values()) / active_threads if active_threads > 0 else 0 - avg_seek_time = sum(stats["avg_seek_time_ms"] for stats in all_stats.values()) / active_threads if active_threads > 0 else 0 - avg_read_time = sum(stats["avg_read_time_ms"] for stats in all_stats.values()) / active_threads if active_threads > 0 else 0 - avg_write_time = sum(stats["avg_write_time_ms"] for stats in all_stats.values()) / active_threads if active_threads > 0 else 0 - - # CORRECTED: Calculate bandwidth averages using corrected formulas - avg_read_bandwidth = sum(stats["avg_read_bandwidth_mib_s"] for stats in all_stats.values()) / active_threads if active_threads > 0 else 0 - avg_write_bandwidth = sum(stats["avg_write_bandwidth_mib_s"] for stats in all_stats.values()) / active_threads if active_threads > 0 else 0 - avg_combined_bandwidth = sum(stats["avg_combined_bandwidth_mib_s"] for stats in all_stats.values()) / active_threads if active_threads > 0 else 0 - - # Calculate combined system bandwidth - combined_read_bandwidth = sum(stats["avg_read_bandwidth_mib_s"] for stats in all_stats.values()) - combined_write_bandwidth = sum(stats["avg_write_bandwidth_mib_s"] for stats in all_stats.values()) - combined_system_bandwidth = sum(stats["avg_combined_bandwidth_mib_s"] for stats in all_stats.values()) - - # Find performance leaders - fastest_thread = max(all_stats.items(), key=lambda x: x[1]["writes_per_second"]) if all_stats else (None, {}) - slowest_thread = min(all_stats.items(), key=lambda x: x[1]["writes_per_second"]) if all_stats else (None, {}) + # Prepare + items = list(stats_map.items()) # [(thread_id, stats), …] + n = len(items) + vals = [s for _, s in items] + + # Sum up all desired metrics in one go + keys = [ + "records_written", "warc_files_created", "total_bytes", "write_errors", + "writes_per_second", "avg_seek_time_ms", "avg_read_time_ms", + "avg_write_time_ms", "avg_read_bandwidth_mib_s", + "avg_write_bandwidth_mib_s", "avg_bandwidth_mib_s" + ] + sums = {k: sum(s[k] for s in vals) for k in keys} + + total_wps = sums["writes_per_second"] + avg = lambda k: sums[k] / n + + # Find fastest and slowest by writes_per_second + fastest_id, fastest = max(items, key=lambda it: it[1]["writes_per_second"]) + slowest_id, slowest = min(items, key=lambda it: it[1]["writes_per_second"]) return { - "active_threads": active_threads, - "total_records_written": total_records, - "total_warc_files_created": total_files, - "total_bytes": total_bytes, - "total_write_errors": total_errors, - "combined_writes_per_second": sum(stats["writes_per_second"] for stats in all_stats.values()), - "avg_writes_per_second_per_thread": avg_writes_per_second, - "avg_seek_time_ms": avg_seek_time, - "avg_read_time_ms": avg_read_time, - "avg_write_time_ms": avg_write_time, - # CORRECTED: Bandwidth metrics using proper calculations - "avg_read_bandwidth_mib_s": avg_read_bandwidth, - "avg_write_bandwidth_mib_s": avg_write_bandwidth, - "avg_combined_bandwidth_mib_s": avg_combined_bandwidth, - "combined_read_bandwidth_mib_s": combined_read_bandwidth, - "combined_write_bandwidth_mib_s": combined_write_bandwidth, - "combined_system_bandwidth_mib_s": combined_system_bandwidth, - "fastest_thread_id": fastest_thread[0], - "fastest_thread_wps": fastest_thread[1].get("writes_per_second", 0), - "slowest_thread_id": slowest_thread[0], - "slowest_thread_wps": slowest_thread[1].get("writes_per_second", 0), - "performance_variance": fastest_thread[1].get("writes_per_second", 0) - slowest_thread[1].get("writes_per_second", 0), - "overall_error_rate_percent": (total_errors / max(total_records + total_errors, 1)) * 100, - "parallelization_efficiency": active_threads * avg_writes_per_second / max(sum(stats["writes_per_second"] for stats in all_stats.values()), 1) if all_stats else 0 + "active_threads": n, + "total_records_written": sums["records_written"], + "total_warc_files_created": sums["warc_files_created"], + "total_bytes": sums["total_bytes"], + "total_write_errors": sums["write_errors"], + "combined_writes_per_second": total_wps, + "avg_writes_per_second_per_thread": avg("writes_per_second"), + "avg_seek_time_ms": avg("avg_seek_time_ms"), + "avg_read_time_ms": avg("avg_read_time_ms"), + "avg_write_time_ms": avg("avg_write_time_ms"), + "avg_read_bandwidth_mib_s": avg("avg_read_bandwidth_mib_s"), + "avg_write_bandwidth_mib_s": avg("avg_write_bandwidth_mib_s"), + "avg_combined_bandwidth_mib_s": avg("avg_bandwidth_mib_s"), + "fastest_thread_id": fastest_id, + "fastest_thread_wps": fastest["writes_per_second"], + "slowest_thread_id": slowest_id, + "slowest_thread_wps": slowest["writes_per_second"], + "performance_variance": fastest["writes_per_second"] - slowest["writes_per_second"], + "overall_error_rate_percent": (sums["write_errors"] / max(sums["records_written"] + sums["write_errors"], 1)) * 100, + "parallelization_efficiency": (n * avg("writes_per_second") / max(total_wps, 1)) } def close_all(self): @@ -1178,7 +1014,7 @@ class HighPerformanceFileProcessor: # Combined bandwidth = 2 * total_bytes / (read_time + write_time) total_io_time = total_read_time + total_write_time - avg_combined_bandwidth_mib_s = (2 * bytes_mib) / max(total_io_time, 0.001) if total_io_time > 0 else 0 + avg_combined_bandwidth_mib_s = (bytes_mib) / max(total_io_time, 0.001) if total_io_time > 0 else 0 return { "success": True, @@ -1330,18 +1166,13 @@ class ZMQStreamingWARCProcessor: self.grouped_tasks = defaultdict(list) self.last_flush_time = time.time() - # Transaction logging - self.transaction_logger = None - if log_path: - transaction_log_path = os.path.join(log_path, "transaction.jsonl") - self.transaction_logger = TransactionLogger(transaction_log_path, enabled=True) - # Parquet job logger with offset-level tracking dest_fs, dest_path = get_fs(config["destination"], warc_location_postfix) self.parquet_logger = ParquetJobLogger( fs=dest_fs, destination_path=str(dest_path), - batch_size=100 + batch_size=100, + refresh_interval_seconds=120 ) # Show resume statistics if in resume mode @@ -1704,6 +1535,8 @@ class ZMQStreamingWARCProcessor: except Exception as e: if self.verbose: self.console.print(f"[red]❌ JobExecutor error: {e}[/red]") + import traceback + traceback.print_exc() time.sleep(0.1) # Wait for remaining futures to complete @@ -1794,13 +1627,6 @@ class ZMQStreamingWARCProcessor: # Update datacenter statistics for failure dc_stats.add_job_failed(processing_time, job.task_count()) - # Log to transaction logger - if self.transaction_logger: - self.transaction_logger.log_file_result( - job.warc_file, job.source_key, job.task_count(), - result["successful"], result["failed"], - result["processing_time"], result.get("error_details", []) - ) def _stats_collector_worker(self): """StatsCollector thread: SUB stats β†’ update metrics (optional).""" @@ -1881,11 +1707,6 @@ class ZMQStreamingWARCProcessor: stats.update({ "running": self.running, - "zmq_architecture": True, - "inproc_transport": True, - "corrected_bandwidth_monitoring": True, # CORRECTED - "expanded_timing_columns": True, - # Destination metrics with corrected bandwidth "destination_stats": destination_stats, "worker_stats": worker_stats, # Detailed per-worker statistics with corrected bandwidth @@ -1898,16 +1719,11 @@ class ZMQStreamingWARCProcessor: "avg_seek_time_ms": destination_stats.get("avg_seek_time_ms", 0), "avg_read_time_ms": destination_stats.get("avg_read_time_ms", 0), "avg_write_time_ms": destination_stats.get("avg_write_time_ms", 0), - # CORRECTED: Bandwidth metrics using proper calculations - "combined_read_bandwidth_mib_s": destination_stats.get("combined_read_bandwidth_mib_s", 0), - "combined_write_bandwidth_mib_s": destination_stats.get("combined_write_bandwidth_mib_s", 0), - "combined_system_bandwidth_mib_s": destination_stats.get("combined_system_bandwidth_mib_s", 0), "avg_read_bandwidth_mib_s": destination_stats.get("avg_read_bandwidth_mib_s", 0), "avg_write_bandwidth_mib_s": destination_stats.get("avg_write_bandwidth_mib_s", 0), "avg_combined_bandwidth_mib_s": destination_stats.get("avg_combined_bandwidth_mib_s", 0), "total_destination_files": destination_stats.get("total_warc_files_created", 0), "total_destination_bytes": destination_stats.get("total_bytes", 0), - # Job log statistics "job_log_stats": job_log_stats, "resume_mode": self.resume_mode, @@ -1922,344 +1738,228 @@ class ZMQStreamingWARCProcessor: return stats def render_status(self) -> Layout: - """Render ZMQ status display with enhanced per-datacenter progress bars, corrected bandwidth statistics, and expanded timing columns.""" status = self.get_status() - - # Compute dynamic height for the progress panel (expanded for more columns) dc_stats = status.get("datacenter_stats", {}) - num_rows = len(dc_stats) * 2 - header_rows = 1 # table header - border_and_pad = 2 # panel top & bottom borders - progress_height = num_rows + header_rows + border_and_pad + worker_stats = status.get("worker_stats", {}) + + _worker_panel, _n_workers = self._workers_panel(worker_stats) + _dc_panel, _n_queues = self._progress_panel(dc_stats) - # Build main layout with dynamic progress size layout = Layout() layout.split_column( Layout(name="header", size=9), - Layout(name="progress", size=progress_height + 2), # +2 for expanded columns - Layout(name="workers", size=14), # +2 for bandwidth columns - Layout(name="stats", size=10), # +2 for bandwidth stats - ) - - # ─── Header with Corrected Bandwidth Overview ──────────────────────────────────────── - header_text = Text() - agg_alive = status.get("aggregator_alive", False) - exec_alive = status.get("executor_alive", False) - stat_alive = status.get("stats_collector_alive", False) - resume_md = status.get("resume_mode", False) - - header_text.append( - f"πŸš€ {'🟒' if agg_alive else 'πŸ”΄'}{'🟒' if exec_alive else 'πŸ”΄'}{'🟒' if stat_alive else 'πŸ”΄'}" - f"{' πŸ”„' if resume_md else ''} ZMQ PUSH/PULL WARC Processor with Corrected Bandwidth Monitoring\n", - style="bold green" + Layout(name="progress", size=_n_queues+7), + Layout(name="workers", size=5+_n_workers//2+1), ) - # CORRECTED: Overall bandwidth summary in header - combined_system_bw = status.get('combined_system_bandwidth_mib_s', 0) - combined_read_bw = status.get('combined_read_bandwidth_mib_s', 0) - combined_write_bw = status.get('combined_write_bandwidth_mib_s', 0) - - if combined_system_bw > 0 or combined_read_bw > 0 or combined_write_bw > 0: - header_text.append( - f"πŸ“‘ Combined Bandwidth: Read {combined_read_bw:.1f} MiB/s, Write {combined_write_bw:.1f} MiB/s, Combined {combined_system_bw:.1f} MiB/s\n", - style="bright_green" - ) - - if resume_md: - js = status.get("job_log_stats", {}) - offs = js.get("completed_offsets", 0) - header_text.append( - f"πŸ”„ Resume Mode (Offset-Level): {js.get('completed', 0)} completed, " - f"{js.get('failed', 0)} failed, {offs} offsets tracked\n", - style="bright_cyan" - ) - - header_text.append( - f"πŸ“Š Query: {status['query_rows_processed']} rows, {status['tasks_sent_to_zmq']} tasks sent\n", - style="cyan" - ) - header_text.append( - f"πŸ”€ ZMQ Pipeline: {status['records_aggregated']} recs aggregated, " - f"{status['jobs_created']} jobs created, {status['avg_tasks_per_job']:.1f} avg tasks/job\n", - style="bright_blue" - ) - header_text.append( - f"πŸ“ˆ Rate: {status['job_rate']:.2f} jobs/s, {status['task_rate']:.2f} tasks/s | " - f"Runtime: {status['elapsed']:.1f}s\n", - style="cyan" - ) - - tsr = status['task_success_rate'] - sc = "green" if tsr > 95 else "yellow" if tsr > 85 else "red" - header_text.append( - f"βœ… Success: {status['tasks_successful']}/{status['tasks_processed']} tasks ({tsr:.1f}%)\n", - style=sc - ) + layout["header"].update(self._header_panel(status)) + layout["progress"].update(_dc_panel) + layout["workers"].update(_worker_panel) + return layout - loss = status['task_loss'] - header_text.append( - "πŸŽ‰ Zero task loss!\n" if loss == 0 else f"⚠️ Task loss: {loss}\n", - style="green" if loss == 0 else "red" - ) - layout["header"].update( - Panel(header_text, title="Status Overview with Corrected Bandwidth", border_style="blue") - ) - # ─── Enhanced Per-Datacenter Progress with Expanded Timing & Corrected Bandwidth ───── - runtime=time.time()-self._start_time + def _progress_panel(self, dc_stats: dict) -> (Panel, int): + if not dc_stats: + return Panel( + "[dim]No datacenter progress available yet[/dim]", + title="Per-Datacenter Progress", border_style="green" + ), 0 + + table = Table(show_header=True, header_style="bold magenta") + # columns + for h, w in [ + ("Datacenter", 10), ("Progress", 16), + ("Q", 8), ("βœ“", 8), ("βœ—", 8), + ("Open [ms]", 12), ("Seek [ms]", 12), + ("Read [ms]", 12), ("Write [ms]", 12), + ("Bandwith [MB/s]", 10), ("Rate", 7), + ("ETA", 6), ("Succ%", 6), + ]: + justify = "center" if h not in ("Datacenter", "Rate", "ETA") else ("right" if h in ("Rate","ETA") else "left") + table.add_column(h, justify=justify, width=w) + + runtime = time.time() - self._start_time total_bytes = 0 - if dc_stats: - progress_table = Table(show_header=True, header_style="bold magenta") - - # Enhanced column layout with expanded timing and corrected bandwidth - progress_table.add_column("Datacenter", style="cyan", width=10) - progress_table.add_column("Progress", width=16) - progress_table.add_column("Q", justify="center", width=8) # Queued (shortened) - progress_table.add_column("βœ“", justify="center", width=8) # Completed - progress_table.add_column("βœ—", justify="center", width=8) # Failed - - # EXPANDED: Four separate timing columns with Β±Οƒ - progress_table.add_column("OpenΒ±Οƒ", justify="center", width=12) - progress_table.add_column("SeekΒ±Οƒ", justify="center", width=12) - progress_table.add_column("ReadΒ±Οƒ", justify="center", width=12) - progress_table.add_column("WriteΒ±Οƒ", justify="center", width=12) - - # CORRECTED: Bandwidth columns - progress_table.add_column("B [MB/s]", justify="center", width=10) # Combined Bandwidth MiB/s - - progress_table.add_column("Rate", justify="right", width=7) - progress_table.add_column("ETA", justify="right", width=6) - progress_table.add_column("Succ%", justify="right", width=6) - - - - for dc, dcst in dc_stats.items(): - total = dcst["total_jobs"] - comp = dcst["jobs_completed"] - fail = dcst["jobs_failed"] - pct = ((comp + fail) / total * 100) if total else 0 - succ_pct = (comp / max(comp + fail, 1) * 100) if total else 0 - - # Compact progress bar - bar_w = 11 # Slightly smaller to fit more columns - filled = int(pct/100 * bar_w) - succ_w = int(comp/total * bar_w) if total else 0 - fail_w = filled - succ_w - empty_w = bar_w - filled - bar = f"[green]{'β–ˆ'*succ_w}[/green][red]{'β–“'*fail_w}[/red][dim]{'β–‘'*empty_w}[/dim]" - pdisp = f"{bar} {pct:.0f}%" if total else "[dim]No jobs[/dim]" - - # Shortened status columns - jobs_s = [f"{dcst['jobs_queued']}", f"{comp}", f"{fail}"] - rate_disp = f"{dcst['completion_rate']:.1f}/s" if dcst["completion_rate"] > 0 else "0/s" - - # ETA - eta = dcst.get("eta_seconds", None) - if eta and eta > 0: - if eta < 60: ed = f"{eta:.0f}s" - elif eta < 3600: ed = f"{eta/60:.0f}m" - else: ed = f"{eta/3600:.1f}h" - else: - ed = "∞" - - # EXPANDED: Individual timing columns with Β±Οƒ (ms) - open_time = (f"{dcst.get('avg_file_open_time_ms', 0):.2g}Β±{dcst.get('file_open_stddev_ms', 0):.0f}" - if dcst.get('avg_file_open_time_ms', 0) > 0 else "[dim]-[/dim]") - seek_time = (f"{dcst.get('avg_seek_time_ms', 0):.2g}Β±{dcst.get('seek_stddev_ms', 0):.0f}" - if dcst.get('avg_seek_time_ms', 0) > 0 else "[dim]-[/dim]") - read_time = (f"{dcst.get('avg_read_time_ms', 0):.2g}Β±{dcst.get('read_stddev_ms', 0):.0f}" - if dcst.get('avg_read_time_ms', 0) > 0 else "[dim]-[/dim]") - write_time = (f"{dcst.get('avg_write_time_ms', 0):.2g}Β±{dcst.get('write_stddev_ms', 0):.0f}" - if dcst.get('avg_write_time_ms', 0) > 0 else "[dim]-[/dim]") - - # CORRECTED: Bandwidth columns (MiB/s) - combined_bw = f"{dcst.get('total_bytes_mb', 0)/runtime:.4g}" - total_bytes += dcst.get('total_bytes_mb', 0) - - # Success percentage with color coding - sc = "green" if succ_pct > 95 else "yellow" if succ_pct > 85 else "red" - success_display = f"[{sc}]{succ_pct:.0f}%[/{sc}]" - - progress_table.add_row( - dc, pdisp, *jobs_s, - open_time, seek_time, read_time, write_time, - combined_bw, - rate_disp, ed, success_display - ) - - layout["progress"].update( - Panel(progress_table, - title="Per-Datacenter Progress with Expanded Timing & Corrected Bandwidth (Open/Seek/Read/Write Β±Οƒ ms, R/W/C BW MiB/s)", - border_style="green") - ) - else: - layout["progress"].update( - Panel("[dim]No datacenter progress available yet[/dim]", - title="Per-Datacenter Progress", border_style="green") - ) - - # ─── Enhanced Worker Performance with Corrected Bandwidth ──────────────────────────── - worker_stats = status.get("worker_stats", {}) - if worker_stats: - wl = list(worker_stats.items()) - mid = len(wl) // 2 - left, right = (wl, []) if len(wl) <= 2 else (wl[:mid], wl[mid:]) # Split differently due to wider table - - wlay = Layout() - if right: - wlay.split_row(Layout(name="workers_left"), Layout(name="workers_right")) - else: - wlay.split_row(Layout(name="workers_left")) - - def make_enhanced_worker_table(entries): - tbl = Table(show_header=True, header_style="bold blue") - - # Enhanced columns with corrected bandwidth - cols = [ - ("Worker", "cyan", 7), - ("Activity", "center", 9), - ("Records", "right", 7), - ("Rate", "right", 7), - ("OpenΒ±Οƒ", "center", 9), - ("SeekΒ±Οƒ", "center", 9), - ("ReadΒ±Οƒ", "center", 9), - ("WriteΒ±Οƒ", "center", 9), - ("R-BW", "center", 6), # CORRECTED: Read Bandwidth - ("W-BW", "center", 6), # CORRECTED: Write Bandwidth - ("C-BW", "center", 6), # CORRECTED: Combined Bandwidth - ("Err%", "right", 5), - ] - - for h, justify, w in cols: - tbl.add_column(h, justify=justify, width=w) - - for tid, st in entries: - act = st["activity_level"].split()[0] - acolor = { - "Very": "bright_green", "Active": "green", - "Moderate": "yellow", "Low": "bright_black", "Idle": "red" - }.get(act, "white") - - # CORRECTED: Bandwidth display - read_bw = f"{st.get('avg_read_bandwidth_mib_s', 0):.1f}" if st.get('avg_read_bandwidth_mib_s', 0) > 0 else "-" - write_bw = f"{st.get('avg_write_bandwidth_mib_s', 0):.1f}" if st.get('avg_write_bandwidth_mib_s', 0) > 0 else "-" - combined_bw = f"{st.get('avg_combined_bandwidth_mib_s', 0):.1f}" if st.get('avg_combined_bandwidth_mib_s', 0) > 0 else "-" - - row = [ - f"T{tid}", - f"[{acolor}]{act}[/{acolor}]", - str(st["records_written"]), - f"{st['writes_per_second']:.1f}/s", - f"{st['avg_file_open_time_ms']:.0f}Β±{st['file_open_stddev_ms']:.0f}", - f"{st['avg_seek_time_ms']:.0f}Β±{st['seek_stddev_ms']:.0f}", - f"{st['avg_read_time_ms']:.0f}Β±{st['read_stddev_ms']:.0f}", - f"{st['avg_write_time_ms']:.0f}Β±{st['write_stddev_ms']:.0f}", - read_bw, # CORRECTED - write_bw, # CORRECTED - combined_bw, # CORRECTED - f"{st['error_rate_percent']:.1f}%" - ] - tbl.add_row(*row) - return tbl - - wlay["workers_left"].update( - Panel(make_enhanced_worker_table(left), - title="List of Threads I", border_style="blue") - ) - if right: - wlay["workers_right"].update( - Panel(make_enhanced_worker_table(right), - title="List of Threads II", border_style="blue") - ) - - layout["workers"].update(wlay) + for dc, d in dc_stats.items(): + total = d["total_jobs"] + comp, fail = d["jobs_completed"], d["jobs_failed"] + pct = (comp + fail) * 100 / total if total else 0 + succ_pct = comp * 100 / max(comp+fail,1) if total else 0 + bar = self._format_bar(comp, fail, total) + eta = self._format_eta(d.get("eta_seconds", 0)) + jobs = [str(d["jobs_queued"]), str(comp), str(fail)] + times = [ + self._format_stat(d, "avg_file_open_time_ms"), + self._format_stat(d, "avg_seek_time_ms"), + self._format_stat(d, "avg_read_time_ms"), + self._format_stat(d, "avg_write_time_ms"), + ] + bw = f"{d.get('total_bytes_mb', 0)/runtime:.4g}" + total_bytes += d.get('total_bytes_mb', 0) + succ_col = self._colorize_pct(succ_pct) + row = [dc, f"{bar} {pct:.0f}%" if total else "[dim]No jobs[/dim]"] + row += jobs + times + [bw, f"{d['completion_rate']:.1f}/s", eta, succ_col] + table.add_row(*row) + + + # --- ADD A TOTAL SUMMARY ROW ACROSS ALL DATACENTERS ---- + total_jobs = sum(d["total_jobs"] for d in dc_stats.values()) + total_queued = sum(d["jobs_queued"] for d in dc_stats.values()) + total_comp = sum(d["jobs_completed"] for d in dc_stats.values()) + total_fail = sum(d["jobs_failed"] for d in dc_stats.values()) + + total_pct = (total_comp + total_fail) * 100 / total_jobs if total_jobs else 0 + total_succ_pct = total_comp * 100 / max(total_comp + total_fail, 1) if total_jobs else 0 + bar = self._format_bar(total_comp, total_fail, total_jobs) + bw_all = f"{total_bytes / runtime:.4g}" + + summary_row = [ + "Total", + f"{bar} {total_pct:.0f}%", + str(total_queued), + str(total_comp), + str(total_fail), + "-", "-", "-", "-", # skip timing columns + bw_all, + "-", # overall rate not shown + "-", # ETA none + self._colorize_pct(total_succ_pct), + ] + table.add_row(*summary_row) + + + return Panel( + table, + title="Per-Datacenter Progress", + border_style="green" + ), len(table.rows) + + def _workers_panel(self, worker_stats: dict) -> (Layout,int): + layout = Layout() + items = list(worker_stats.items()) + mid = len(items) // 2 or 1 + left, right = items[:mid], items[mid:] + if right: + layout.split_row(Layout(name="left"), Layout(name="right")) + layout["left"].update(self._worker_panel(left, "Threads I")) + layout["right"].update(self._worker_panel(right, "Threads II")) else: - layout["workers"].update( - Panel("[dim]No worker performance data available yet[/dim]", - title="Worker Performance", border_style="blue") - ) - - # ─── Enhanced Stats Summary with Corrected Bandwidth ───────────────────────────────── - stats_text = Text() - - # ZMQ errors - se, re = status['zmq_send_errors'], status['zmq_recv_errors'] - if se > 0 or re > 0: - stats_text.append(f"⚠️ ZMQ Errors: {se} send, {re} recv\n", style="yellow") - - # Thread health - agg_done = self.record_aggregator_finished.is_set() - stats_text.append( - f"🧡 Threads: Aggregator {'🏁' if agg_done else ('🟒' if agg_alive else 'πŸ”΄')}, " - f"Executor {'🟒' if exec_alive else 'πŸ”΄'}, " - f"StatsCollector {'🟒' if stat_alive else 'πŸ”΄'}\n", - style="bright_magenta" + layout.split_row(Layout(name="left")) + layout["left"].update(self._worker_panel(left, "Threads")) + return layout, len(items) + + def _worker_panel(self, entries: list, title: str) -> Panel: + tbl = Table(show_header=True, header_style="bold blue") + cols = [ + ("Worker", 8), ("Act", 9), ("Recs", 6), ("Rate", 6), + ("Open", 5), ("Seek", 5), ("Read", 5), ("Write", 5), + ("R-BW", 8), ("W-BW", 8), ("C-BW", 8), ("Err%", 5), + ] + for h, w in cols: + tbl.add_column(h, justify="center", width=w) + + for tid, st in entries: + act = st["activity_level"].split()[0] + color = { + "Very":"bright_green","Active":"green", + "Moderate":"yellow","Low":"bright_black","Idle":"red" + }[act] + times = [ + f"{st['avg_file_open_time_ms']:.0f}", + f"{st['avg_seek_time_ms']:.0f}", + f"{st['avg_read_time_ms']:.0f}", + f"{st['avg_write_time_ms']:.0f}", + ] + bws = [ + f"{st.get('avg_read_bandwidth_mib_s',0):.1f}" or "-", + f"{st.get('avg_write_bandwidth_mib_s',0):.1f}" or "-", + f"{st.get('avg_bandwidth_mib_s',0):.1f}" or "-", + ] + row = [ + f"T{tid}", + f"[{color}]{act}[/{color}]", + str(st["records_written"]), + f"{st['writes_per_second']:.1f}/s", + *times, *bws, + f"{st['error_rate_percent']:.1f}%" + ] + tbl.add_row(*row) + + return Panel(tbl, title=title, border_style="blue") + + def _header_panel(self, status: dict) -> Panel: + t = Text() + flags = "->".join("🟒" if status.get(k) else "πŸ”΄" + for k in ("aggregator_alive","executor_alive","stats_collector_alive")) + t.append(f"πŸš€ {flags}{' πŸ”„' if status.get('resume_mode') else ''} Pipeline Status (Aggregator->Executor->Stats Collector)\n", + style="bold green") + + if status.get("resume_mode"): + js = status.get("job_log_stats", {}) + t.append(f"πŸ”„ Resume: {js.get('completed',0)}βœ“/{js.get('failed',0)}βœ—, " + f"{js.get('completed_offsets',0)} offsets\n", + style="bright_cyan") + + t.append(f"πŸ“Š Flow: {status['query_rows_processed']} rows -> {status['tasks_sent_to_zmq']} tasks (offsets) -> " + f"{status['jobs_created']} jobs (files) with {status['avg_tasks_per_job']:.1f} offsets/file\n", style="cyan") + t.append(f"πŸ“ˆ Rate: {status['job_rate']:.2f} Jobs/s, {status['task_rate']:.2f} Tasks/s | " + f"Runtime: {status['elapsed']:.1f}s\n", style="cyan") + + alive = status['active_threads'] + if alive: + eff = status.get('parallelization_efficiency',0) * 100 + t.append(f"πŸ‘₯ Workers: {alive}, {status.get('combined_writes_per_second',0):.1f} writes/s, " + f"{eff:.1f}% effectiveness\n", style="bright_blue") + + tsr = status.get('task_success_rate',0) + color = "green" if tsr>95 else "yellow" if tsr>85 else "red" + t.append(f"βœ… Success: {status['tasks_successful']}/{status['tasks_processed']} ({tsr:.1f}%)\n", + style=color) + + loss = status.get('task_loss',0) + t.append( + "πŸŽ‰ Zero task loss!\n" if loss==0 else f"⚠️ Task loss: {loss}\n", + style="green" if loss==0 else "red" ) - # CORRECTED: Enhanced bandwidth summary - at = status.get('active_threads', 0) - if at: - cwps = status.get('combined_writes_per_second', 0) - eff = status.get('parallelization_efficiency', 0) * 100 - - # Calculate corrected total bandwidth across all workers - total_read_bw = status.get('combined_read_bandwidth_mib_s', 0) - total_write_bw = status.get('combined_write_bandwidth_mib_s', 0) - total_combined_bw = status.get('combined_system_bandwidth_mib_s', 0) - - stats_text.append( - f"πŸ‘₯ Workers: {at} active, {cwps:.1f} writes/s, {eff:.1f}% efficiency\n" - f"πŸ“‘ Corrected System Bandwidth: R{total_read_bw:.1f} + W{total_write_bw:.1f} = C{total_combined_bw:.1f} MiB/s\n", - style="bright_blue" - ) - - # Bandwidth per worker analysis with corrected calculations - if worker_stats: - worker_combined_bandwidths = [ - w.get('avg_combined_bandwidth_mib_s', 0) for w in worker_stats.values() - ] - if worker_combined_bandwidths: - avg_combined_bw_per_worker = sum(worker_combined_bandwidths) / len(worker_combined_bandwidths) - max_combined_bw = max(worker_combined_bandwidths) - min_combined_bw = min(worker_combined_bandwidths) - stats_text.append( - f" Combined BW Range: {min_combined_bw:.1f} - {max_combined_bw:.1f} MiB/s per worker (avg: {avg_combined_bw_per_worker:.1f})\n", - style="bright_cyan" - ) + se, re = status.get('zmq_send_errors',0), status.get('zmq_recv_errors',0) + if se or re: + t.append(f"⚠️ ZMQ Errors: {se} send, {re} recv\n", style="yellow") - # File grouping with bandwidth efficiency - fg = status.get('files_grouped', 0) - if fg: - stats_text.append( - f"πŸ“¦ File Grouping: {fg} files grouped, {status['avg_tasks_per_job']:.1f} avg tasks/job\n", - style="bright_green" - ) + return Panel(t, title="Status Overview", border_style="blue") - # Overall DC timing with corrected bandwidth - if dc_stats: - total_c = sum(d["jobs_completed"] for d in dc_stats.values()) - if total_c: - ws = sum(d.get("avg_seek_time_ms", 0) * d["jobs_completed"] for d in dc_stats.values()) / total_c - wr = sum(d.get("avg_read_time_ms", 0) * d["jobs_completed"] for d in dc_stats.values()) / total_c - ww = sum(d.get("avg_write_time_ms", 0) * d["jobs_completed"] for d in dc_stats.values()) / total_c - - # CORRECTED: Overall bandwidth calculation - total_read_bw_dc = sum(d.get("avg_read_bandwidth_mib_s", 0) for d in dc_stats.values()) - total_write_bw_dc = sum(d.get("avg_write_bandwidth_mib_s", 0) for d in dc_stats.values()) - total_combined_bw_dc = sum(d.get("avg_combined_bandwidth_mib_s", 0) for d in dc_stats.values()) - - stats_text.append( - f"⚑ Overall Performance: Seek {ws:.1f}ms, Read {wr:.1f}ms, Write {ww:.1f}ms\n" - f"πŸ“Š Overall Corrected Bandwidth: Read {total_read_bw_dc:.1f} MiB/s, Write {total_write_bw_dc:.1f} MiB/s, Combined {total_bytes/runtime:.4g} MiB/s, {total_bytes:.4g} MB, Runtime: {runtime}\n", - style="yellow" - ) + # ─── Utility formatting ──────────────────────────────────────────────────────────────── - layout["stats"].update( - Panel(stats_text, title="Enhanced Performance & Corrected Bandwidth Architecture", border_style="yellow") + def _format_bar(self, comp, fail, total, width: int = 11) -> str: + filled = int((comp+fail)/total * width) if total else 0 + succ = int(comp/total * width) if total else 0 + fail_w = filled - succ + empty = width - filled + return ( + f"[green]{'β–ˆ'*succ}[/green]" + f"[red]{'β–“'*fail_w}[/red]" + f"[dim]{'β–‘'*empty}[/dim]" ) - return layout + def _format_eta(self, seconds: float) -> str: + if not seconds or seconds <= 0: + return "∞" + if seconds < 60: + return f"{seconds:.0f}s" + if seconds < 3600: + return f"{seconds/60:.0f}m" + return f"{seconds/3600:.1f}h" + + def _format_stat(self, d: dict, avg_key: str) -> str: + avg = d.get(avg_key, 0) + if avg <= 0: + return "[dim]-[/dim]" + return f"{avg:.2g}" + + def _colorize_pct(self, pct: float) -> str: + c = "green" if pct>95 else "yellow" if pct>85 else "red" + return f"[{c}]{pct:.0f}%[/{c}]" def start(self): """Start the ZMQ streaming processor.""" @@ -2290,7 +1990,7 @@ class ZMQStreamingWARCProcessor: self.console.print(f"[green]πŸš€ Started ZMQ PUSH/PULL processor with corrected bandwidth monitoring: {self.max_workers} workers{resume_info}[/green]") def stop(self, timeout: float = 60.0): - """Stop the ZMQ streaming processor.""" + """Stop the ZMQ streaming processor and generate completion report.""" if not self.running: return @@ -2328,108 +2028,236 @@ class ZMQStreamingWARCProcessor: # Close destination instances self.destination_manager.close_all() + # Generate comprehensive markdown report + try: + completion_report_md = self.generate_completion_report_markdown() + + # Store the report via parquet logger + self.parquet_logger.store_completion_report(completion_report_md) + + if self.verbose: + self.console.print("[cyan]πŸ“Š Completion report generated and stored[/cyan]") + + except Exception as e: + if self.verbose: + self.console.print(f"[yellow]⚠️ Could not generate completion report: {e}[/yellow]") + # Flush log entries self.parquet_logger.flush() # Cleanup ZMQ self._cleanup_zmq() - # Print completion statistics with corrected bandwidth - job_stats = self.parquet_logger.get_completion_stats() - performance_stats = self.parquet_logger.get_performance_stats() + # Print basic completion message (detailed report will be printed by warc() function) final_stats = self.get_status() + self.console.print(f"[green]βœ… ZMQ streaming processor stopped[/green]") + self.console.print(f"[green]πŸ“Š Final: {final_stats['jobs_completed']} jobs, {final_stats['records_written']} records, " + f"success rate: {final_stats['task_success_rate']:.1f}%[/green]") - self.console.print(f"[cyan]πŸ“Š Job completion statistics: {job_stats}[/cyan]") + # Zero task loss verification + task_loss = final_stats.get('task_loss', 0) + if task_loss == 0: + self.console.print("[green]πŸŽ‰ Zero task loss achieved![/green]") + else: + self.console.print(f"[red]⚠️ {task_loss} tasks lost[/red]") - if performance_stats: - self.console.print(f"[cyan]⚑ Performance Summary: " - f"Seek {performance_stats.get('avg_seek_time_ms', 0):.1f}ms, " - f"Read {performance_stats.get('avg_read_time_ms', 0):.1f}ms, " - f"Write {performance_stats.get('avg_write_time_ms', 0):.1f}ms | " - f"{performance_stats.get('bytes_per_second', 0)/1024/1024:.1f} MB/s[/cyan]") + def generate_completion_report_markdown(self) -> str: + """Generate comprehensive completion report in Markdown format.""" + final_stats = self.get_status() + job_stats = self.parquet_logger.get_completion_stats() + performance_stats = self.parquet_logger.get_performance_stats() + + # Calculate runtime + total_time = time.time() - self._start_time - # Per-datacenter completion summary with corrected bandwidth + # Get key metrics dc_stats = final_stats.get("datacenter_stats", {}) + worker_stats = final_stats.get("worker_stats", {}) + + # ZMQ error metrics + zmq_send_errors = final_stats.get('zmq_send_errors', 0) + zmq_recv_errors = final_stats.get('zmq_recv_errors', 0) + total_zmq_messages = final_stats.get('tasks_sent_to_zmq', 0) + final_stats.get('jobs_created', 0) + total_bytes_mb = sum([v.get("total_bytes_mb",0) for k,v in final_stats['datacenter_stats'].items()]) + + # Thread health + aggregator_alive = final_stats.get('aggregator_alive', False) + executor_alive = final_stats.get('executor_alive', False) + stats_collector_alive = final_stats.get('stats_collector_alive', False) + healthy_threads = sum([aggregator_alive, executor_alive, stats_collector_alive]) + + # Calculate assessment metrics + task_loss = final_stats.get('task_loss', 0) + success_rate = final_stats.get('task_success_rate', 0) + zmq_message_reliability = (zmq_send_errors + zmq_recv_errors) == 0 + + # Worker performance assessment + worker_performance_good = True + bandwidth_performance_good = True + if worker_stats: + worker_rates = [w['writes_per_second'] for w in worker_stats.values()] + worker_combined_bandwidths = [w.get('avg_combined_bandwidth_mib_s', 0) for w in worker_stats.values()] + + if worker_rates: + rate_stddev = statistics.stdev(worker_rates) if len(worker_rates) > 1 else 0.0 + avg_rate = statistics.mean(worker_rates) + coefficient_of_variation = (rate_stddev / max(avg_rate, 1)) * 100 + worker_performance_good = coefficient_of_variation < 50 + + if worker_combined_bandwidths: + bw_stddev = statistics.stdev(worker_combined_bandwidths) if len(worker_combined_bandwidths) > 1 else 0.0 + avg_bw = statistics.mean(worker_combined_bandwidths) + bw_coefficient_of_variation = (bw_stddev / max(avg_bw, 1)) * 100 if avg_bw > 0 else 0 + bandwidth_performance_good = bw_coefficient_of_variation < 50 + + # Start building the markdown report + md = [] + + # Header + md.append("# WARC Cache Fetch Report") + md.append("") + md.append(f"\n**Processing Date:** {datetime.now().isoformat()}") + md.append(f"\n**Total Runtime:** {total_time:.1f} seconds") + md.append(f"\n**Total MiB:** {total_bytes_mb:.2f} ") + md.append(f"\n**Resume Mode:** {'Enabled' if self.resume_mode else 'Disabled'}") + + # Key Metrics + md.append("## Key Metrics") + md.append("") + md.append("| Metric | Value |") + md.append("|--------|-------|") + md.append(f"| Jobs Completed | {final_stats['jobs_completed']:,} |") + md.append(f"| Records Written | {final_stats['records_written']:,} |") + md.append(f"| Success Rate | {success_rate:.1f}% |") + md.append(f"| Task Loss | {task_loss:,} |") + md.append(f"| Processing Rate | {final_stats['records_written'] / max(total_time, 1):.1f} records/second |") + md.append(f"| Bandwidth | {total_bytes_mb/total_time:.3f} MiB/s |") + md.append("") + + # ZMQ Architecture Summary + md.append("## Queuing Summary") + md.append("") + md.append("| Component | Details |") + md.append("|-----------|---------|") + md.append(f"| High-water Mark | {self.zmq_hwm} (automatic back-pressure) |") + md.append(f"| Records Sent to ZMQ | {final_stats.get('tasks_sent_to_zmq', 0):,} |") + md.append(f"| Records Aggregated | {final_stats.get('records_aggregated', 0):,} |") + md.append(f"| Jobs Created | {final_stats.get('jobs_created', 0):,} |") + md.append(f"| Files Grouped | {final_stats.get('files_grouped', 0):,} |") + md.append(f"| Avg Tasks per Job | {final_stats.get('avg_tasks_per_job', 0):.1f} |") + md.append("") + + # Per-Datacenter Performance if dc_stats: - self.console.print(f"[cyan]🏒 Per-Datacenter Completion Summary with Corrected Bandwidth:[/cyan]") + md.append("## Per-Datacenter Performance Analysis") + md.append("") + + total_dc_jobs = sum(dc["total_jobs"] for dc in dc_stats.values()) + total_dc_completed = sum(dc["jobs_completed"] for dc in dc_stats.values()) + total_dc_failed = sum(dc["jobs_failed"] for dc in dc_stats.values()) + total_read_bw, total_write_bw, total_bytes = 0, 0, 0 + + md.append("| Datacenter | Jobs | Success Rate | Performance (O/S/R/W ms) | Bandwidth (R/W/C MiB/s) |") + md.append("|------------|------|-------------|--------------------------|-------------------------|") + for dc_name, dc_stat in dc_stats.items(): - total_jobs = dc_stat["total_jobs"] - completed = dc_stat["jobs_completed"] - failed = dc_stat["jobs_failed"] - success_rate = dc_stat["success_rate"] + dc_pct_of_total = (dc_stat["total_jobs"] / max(total_dc_jobs, 1)) * 100 + + # Performance metrics + timing_summary = f"{dc_stat.get('avg_file_open_time_ms', 0):.1f}/{dc_stat.get('avg_seek_time_ms', 0):.1f}/{dc_stat.get('avg_read_time_ms', 0):.1f}/{dc_stat.get('avg_write_time_ms', 0):.1f}" + + # Bandwidth metrics read_bw = dc_stat.get("avg_read_bandwidth_mib_s", 0) write_bw = dc_stat.get("avg_write_bandwidth_mib_s", 0) - combined_bw = dc_stat.get("avg_combined_bandwidth_mib_s", 0) + combined_bw = dc_stat.get("total_bytes_mb", 0) / dc_stat.get("elapsed", 0) if dc_stat.get("elapsed", 0) > 0 else 0 + bandwidth_summary = f"{read_bw:.1f}/{write_bw:.1f}/{combined_bw:.1f}" + total_read_bw += read_bw + total_write_bw += write_bw + total_bytes += dc_stat.get("total_bytes_mb", 0) + job_summary = f"{dc_stat['jobs_completed']}/{dc_stat['total_jobs']} ({dc_stat['jobs_failed']} failed)" - self.console.print(f" β€’ {dc_name}: {completed}/{total_jobs} jobs completed " - f"({failed} failed, {success_rate:.1f}% success rate) | " - f"BW: R{read_bw:.1f} W{write_bw:.1f} C{combined_bw:.1f} MiB/s") + md.append(f"| {dc_name} | {job_summary} | {dc_stat['success_rate']:.1f}% | {timing_summary} | {bandwidth_summary} |") - # Worker performance summary with corrected bandwidth - worker_stats = final_stats.get("worker_stats", {}) - if worker_stats: - active_workers = len(worker_stats) - total_worker_records = sum(stats["records_written"] for stats in worker_stats.values()) - avg_worker_rate = sum(stats["writes_per_second"] for stats in worker_stats.values()) / max(active_workers, 1) - total_worker_read_bw = sum(stats.get("avg_read_bandwidth_mib_s", 0) for stats in worker_stats.values()) - total_worker_write_bw = sum(stats.get("avg_write_bandwidth_mib_s", 0) for stats in worker_stats.values()) - total_worker_combined_bw = sum(stats.get("avg_combined_bandwidth_mib_s", 0) for stats in worker_stats.values()) - - self.console.print(f"[cyan]πŸ‘₯ Worker Performance Summary with Corrected Bandwidth:[/cyan]") - self.console.print(f" β€’ Active Workers: {active_workers}") - self.console.print(f" β€’ Total Records: {total_worker_records:,}") - self.console.print(f" β€’ Average Rate: {avg_worker_rate:.1f} writes/s per worker") - self.console.print(f" β€’ System Bandwidth: Read {total_worker_read_bw:.1f} MiB/s, Write {total_worker_write_bw:.1f} MiB/s, Combined {total_worker_combined_bw:.1f} MiB/s") - - # Show top performers with corrected bandwidth - sorted_workers = sorted(worker_stats.items(), key=lambda x: x[1]["writes_per_second"], reverse=True) - if len(sorted_workers) >= 2: - best = sorted_workers[0] - worst = sorted_workers[-1] - best_combined_bw = best[1].get("avg_combined_bandwidth_mib_s", 0) - worst_combined_bw = worst[1].get("avg_combined_bandwidth_mib_s", 0) - self.console.print(f" β€’ Best Performer: Thread {best[0]} ({best[1]['writes_per_second']:.1f}/s, {best_combined_bw:.1f} MiB/s)") - self.console.print(f" β€’ Needs Attention: Thread {worst[0]} ({worst[1]['writes_per_second']:.1f}/s, {worst_combined_bw:.1f} MiB/s)") - - # ZMQ specific statistics - stats = self.metrics.get_stats() - self.console.print(f"[cyan]πŸ”Œ ZMQ Statistics:[/cyan]") - self.console.print(f" β€’ Records sent to ZMQ: {stats['tasks_sent_to_zmq']}") - self.console.print(f" β€’ Records aggregated: {stats['records_aggregated']}") - self.console.print(f" β€’ Jobs created: {stats['jobs_created']}") - self.console.print(f" β€’ Average tasks per job: {stats['avg_tasks_per_job']:.1f}") - self.console.print(f" β€’ ZMQ send errors: {stats['zmq_send_errors']}") - self.console.print(f" β€’ ZMQ recv errors: {stats['zmq_recv_errors']}") - - # Final statistics with corrected bandwidth - destination_stats = self.destination_manager.get_aggregated_stats() + bandwidth_summary = f"{total_read_bw:.1f}/{total_write_bw:.1f}/{total_bytes_mb/total_time:.3f}" - self.console.print(f"[green]βœ… ZMQ streaming processor stopped[/green]") - self.console.print(f"[green]πŸ“Š Final: {stats['jobs_completed']} jobs, {stats['records_written']} records, " - f"success rate: {stats['task_success_rate']:.1f}%[/green]") - - # Enhanced parallel performance summary with corrected bandwidth - if destination_stats: - active_threads = destination_stats.get('active_threads', 0) - combined_wps = destination_stats.get('combined_writes_per_second', 0) - efficiency = destination_stats.get('parallelization_efficiency', 0) - combined_read_bw = destination_stats.get('combined_read_bandwidth_mib_s', 0) - combined_write_bw = destination_stats.get('combined_write_bandwidth_mib_s', 0) - combined_system_bw = destination_stats.get('combined_system_bandwidth_mib_s', 0) - - self.console.print(f"[green]πŸš€ Enhanced Parallel Performance: {active_threads} threads, " - f"{combined_wps:.1f} combined writes/s, " - f"{efficiency*100:.1f}% efficiency | " - f"Corrected Bandwidth: R{combined_read_bw:.1f} W{combined_write_bw:.1f} C{combined_system_bw:.1f} MiB/s[/green]") + overall_dc_success = (total_dc_completed / max(total_dc_completed + total_dc_failed, 1)) * 100 + md.append(f"| **TOTAL** | **{total_dc_completed}/{total_dc_jobs}** | **{overall_dc_success:.1f}%** | - | {bandwidth_summary} |") + md.append("") + + # Job Completion Statistics + md.append("## Job Completion Statistics") + md.append("") + md.append("**Job Log Summary:**") + md.append("") + md.append("| Key | Value |") + md.append("| --- | ----- |") + for key, value in job_stats.items(): + md.append(f"| {key} | {value} |") + + md.append("") + + if self.resume_mode: + md.append("### Resume Mode Statistics") + md.append("") + js = final_stats.get("job_log_stats", {}) + md.append(f"- **Previously Completed:** {js.get('completed',0)} jobs") + md.append(f"- **Previously Failed:** {js.get('failed',0)} jobs") + md.append(f"- **Completed Offsets:** {js.get('completed_offsets',0)}") + md.append("") + + # Configuration + md.append("## Configuration") + md.append("") + md.append("| Parameter | Value |") + md.append("|-----------|-------|") + md.append(f"| Max Workers | {self.max_workers} |") + md.append(f"| Record Threshold | {self.record_threshold} |") + md.append(f"| Time Threshold | {self.time_threshold}s |") + md.append(f"| ZMQ HWM | {self.zmq_hwm} |") + md.append(f"| Stats Interval | {self.stats_interval}s |") + md.append("") + + return "\n".join(md) + + def print_completion_report(self) -> None: + """Print the stored completion report using Rich UI.""" + try: + # Try to get stored markdown report from parquet logger + report_md = self.parquet_logger.get_completion_report_markdown() + if not report_md: + # Fallback to generating a new report + report_md = self.generate_completion_report_markdown() + + # Parse and render with Rich + from rich.markdown import Markdown + from rich.panel import Panel + + markdown_obj = Markdown(report_md) + + # Wrap in a panel for better presentation + panel = Panel( + markdown_obj, + title="πŸŽ‰ Processing Complete - Final Report", + border_style="green", + expand=False + ) + + self.console.print("\n") + self.console.print(panel) + + except Exception as e: + # Fallback to basic text output + self.console.print(f"[red]❌ Error rendering completion report: {e}[/red]") + self.console.print("[yellow]πŸ“„ Generating basic completion summary...[/yellow]") + + final_stats = self.get_status() + success_rate = final_stats.get('task_success_rate', 0) + jobs_completed = final_stats.get('jobs_completed', 0) + records_written = final_stats.get('records_written', 0) + + self.console.print(f"[green]βœ… Processing completed: {jobs_completed} jobs, {records_written} records, {success_rate:.1f}% success rate[/green]") - # Zero task loss verification - if stats['task_loss'] == 0: - self.console.print("[green]πŸŽ‰ Zero task loss achieved![/green]") - else: - self.console.print(f"[red]⚠️ {stats['task_loss']} tasks lost[/red]") - # Enhanced ZMQ architecture summary - self.console.print("[green]πŸ”Œ ZMQ PUSH/PULL Architecture with Corrected Bandwidth Monitoring: Reliable message delivery + automatic back-pressure + accurate bandwidth analysis achieved[/green]") def __enter__(self): """Context manager entry.""" @@ -2442,43 +2270,31 @@ class ZMQStreamingWARCProcessor: def warc(self, local_specifier: str, remote_specifier: str, urls_file: str = None, where: Optional[str] = "", limit: Optional[int] = None, files: str = "**/*.parquet", - pq_batch_size: int = 1, batch_size: int = 100, prefetch: int = 10, page_size: int = 10, + pq_batch_size: int = 1, batch_size: int = 1000, prefetch: int = 100, page_size: int = 10, rollover_limit: int = 3000, max_workers: int = 10, warc_location_cfg: Optional[str] = None, verbose: bool = False, log_path: str = "", record_threshold: int = 1000, time_threshold: float = 30.0, - zmq_hwm: int = 1000, resume: bool = False, stats_interval: float = 5.0, + zmq_hwm: int = 1000, resume: bool = False, stats_interval: float = 8.0, group_name:str =''): """ - High-performance parallel WARC processing with ZeroMQ PUSH/PULL architecture, - detailed per-datacenter progress tracking, comprehensive per-worker performance monitoring, - corrected real-time bandwidth analysis (MiB/s), and expanded timing display columns. + WARC Cache fetch function. for a given selected dataset the query is executed and the WARC files for the found URLs are fetched. + The system utilizes message queues and threading in order to optimize bandwith. However, note that if teh number of + records is too large, fetching full warc files might be better. - This enhanced processor implements a ZeroMQ-based task distribution system: - - Main Thread: SQL streaming β†’ PUSH WARCTask to RecordAggregator - - RecordAggregator: PULL WARCTask β†’ group by file β†’ PUSH FileJob to JobExecutor - - JobExecutor: PULL FileJob β†’ submit to ThreadPoolExecutor for I/O processing - - StatsCollector: SUB stats β†’ update metrics (optional telemetry) + For optimal performance we also recommend to query only local datastes (i.e. use owilix remote pull before). - Key improvements over queue-based architecture: - - Automatic back-pressure via ZMQ high-water marks (HWM) - - Reliable message delivery through PUSH/PULL pattern - - Non-blocking telemetry via PUB/SUB - - Future-proof scaling with transport flexibility (inproc β†’ ipc β†’ tcp) - - Maintained compatibility with existing logging and resume functionality - - Per-datacenter progress tracking with Rich UI progress bars - - Intelligent ETA calculation accounting for growing queues - - Comprehensive per-worker performance monitoring with timing analysis and standard deviations - - Enhanced datacenter statistics with seek/read/write performance metrics and variance analysis - - CORRECTED: Real-time bandwidth monitoring (MiB/s) with proper statistical variance tracking - - CORRECTED: Expanded timing display with four separate columns (Open/Seek/Read/Write Β±Οƒ ms) - - CORRECTED: Performance consistency analysis across workers and datacenters + Accessing the WARC Cache requires proper credentials and configuration of the warc cache. please refer to the documentation. - Bandwidth Calculation (CORRECTED): - 1. Single total_bytes tracking (no distinction between read/write bytes) - 2. Read bandwidth = total_bytes / read_time (MiB/s) - 3. Write bandwidth = total_bytes / write_time (MiB/s) - 4. Combined bandwidth = 2 * total_bytes / (read_time + write_time) (MiB/s) - 5. Statistical variance analysis across workers and datacenters + The process displays a sophisticated console ui with detailed statistic. + All fetches are logged in a parquet transaction log, such that operations can be resumed. + + Usage Example: + + ``` + python -m owilix.cli query warc --local "all:2025-06-06#1/collectionName=main" \ + where="ows_genai IS NOT NULL AND ows_genai=TRUE" \ + warc_location_cfg=.env-warc-cfg.json urls_file=/home/mgrani/tmp/oem.csv max_workers=20 record_threshold=20000 zmq_hwm=100000 resume=True group_name=all_250606-2D + ``` Args: local_specifier (str): Local filesystem path pattern for input files @@ -2520,10 +2336,11 @@ def warc(self, local_specifier: str, remote_specifier: str, urls_file: str = Non total_processed = 0 all_final_stats = {} + processor = None # Initialize processor variable # Process each group with ZMQ streaming processor if True: # hack for removing groups - self.console.print(f"\n[bold cyan]πŸ“‹ Processing group {group_name} with Enhanced ZMQ + Corrected Bandwidth Monitoring of {sum([len(f) for f in all_files.values()])}[/bold cyan]") + self.console.print(f"\n[bold cyan]πŸ“‹ Processing group {group_name} with {sum([len(f) for f in all_files.values()])} files [/bold cyan]") # Use Enhanced ZMQ Streaming WARC Processor with Per-Datacenter Progress, Worker Performance, and Corrected Bandwidth Monitoring with ZMQStreamingWARCProcessor( @@ -2597,7 +2414,7 @@ def warc(self, local_specifier: str, remote_specifier: str, urls_file: str = Non return processed_rows # Start Enhanced ZMQ processor with corrected bandwidth monitoring - self.console.print("[cyan]🌊 Starting Enhanced ZMQ PUSH/PULL processing with per-datacenter, per-worker, and corrected bandwidth monitoring...[/cyan]") + self.console.print("[cyan]🌊 Starting WARC Cache fetching[/cyan]") processor.start() # Process query results with Enhanced Live display including corrected bandwidth @@ -2608,17 +2425,21 @@ def warc(self, local_specifier: str, remote_specifier: str, urls_file: str = Non try: return processor.render_status() except Exception: + from rich.text import Text text = Text() - text.append("πŸš€ Enhanced ZMQ PUSH/PULL WARC Processor with Corrected Bandwidth Monitoring\n", style="bold green") text.append("Status rendering temporarily unavailable\n", style="yellow") + if self.verbose: + import traceback + traceback.print_exc() + return text # Process query results with Enhanced ZMQ streaming processor self._process_query_results( db, all_files, False, page_size, fn_callback=enhanced_zmq_streaming_result_processor, - task_str=f"streaming records for group {group_name} (Enhanced ZMQ PUSH/PULL with Per-DC, Worker Performance & Corrected Bandwidth Monitoring)", + task_str=f"streaming records for group {group_name}", external_live=live, external_content_getter=update_display ) @@ -2776,8 +2597,9 @@ def warc(self, local_specifier: str, remote_specifier: str, urls_file: str = Non pass except Exception as e: - self.console.print(f"[red]❌ Enhanced ZMQ live display error: {e}[/red]") - + self.console.print(f"[red]❌ Display error: {e}[/red]") + import traceback + traceback.print_exc() # Fallback processing without live display self._process_query_results( db, all_files, False, page_size, @@ -2785,338 +2607,38 @@ def warc(self, local_specifier: str, remote_specifier: str, urls_file: str = Non task_str=f"streaming records for group {group_name} (Enhanced ZMQ fallback mode)" ) - self.console.print("[cyan]βœ… Enhanced ZMQ processing complete, stopping processor...[/cyan]") + self.console.print("[cyan]βœ… Processing complete, stopping processor...[/cyan]") # Get final Enhanced ZMQ stats with per-datacenter, worker details, and corrected bandwidth group_stats = processor.get_status() total_processed += group_stats["jobs_completed"] all_final_stats = group_stats - # Enhanced ZMQ group completion summary with per-datacenter, worker breakdown, and corrected bandwidth - success_rate = group_stats.get('task_success_rate', 0) - dc_stats = group_stats.get("datacenter_stats", {}) - worker_stats = group_stats.get("worker_stats", {}) - combined_read_bw = group_stats.get('combined_read_bandwidth_mib_s', 0) - combined_write_bw = group_stats.get('combined_write_bandwidth_mib_s', 0) - combined_system_bw = group_stats.get('combined_system_bandwidth_mib_s', 0) - - self.console.print(f"[green]βœ… Group {group_name} completed with enhanced monitoring: " - f"{group_stats['jobs_completed']} jobs processed, " - f"{group_stats['records_written']} records written, " - f"success rate: {success_rate:.1f}% | " - f"Corrected Bandwidth: R{combined_read_bw:.1f} W{combined_write_bw:.1f} C{combined_system_bw:.1f} MiB/s[/green]") - - # Enhanced per-datacenter completion summary with corrected bandwidth - if dc_stats: - self.console.print(f"[bright_green]🏒 Enhanced Per-Datacenter Completion Summary with Corrected Bandwidth:[/bright_green]") - for dc_name, dc_stat in dc_stats.items(): - total_dc_jobs = dc_stat["total_jobs"] - completed_dc_jobs = dc_stat["jobs_completed"] - failed_dc_jobs = dc_stat["jobs_failed"] - dc_success_rate = dc_stat["success_rate"] - read_bw = dc_stat.get("avg_read_bandwidth_mib_s", 0) - write_bw = dc_stat.get("avg_write_bandwidth_mib_s", 0) - combined_bw = dc_stat.get("avg_combined_bandwidth_mib_s", 0) - - # Enhanced timing summary for datacenter with corrected bandwidth - timing_summary = "" - if "avg_seek_time_ms" in dc_stat and dc_stat["avg_seek_time_ms"] > 0: - timing_summary = (f" | Performance: Open {dc_stat.get('avg_file_open_time_ms', 0):.1f}Β±{dc_stat.get('file_open_stddev_ms', 0):.1f}ms, " - f"Seek {dc_stat['avg_seek_time_ms']:.1f}Β±{dc_stat.get('seek_stddev_ms', 0):.1f}ms, " - f"Read {dc_stat['avg_read_time_ms']:.1f}Β±{dc_stat.get('read_stddev_ms', 0):.1f}ms, " - f"Write {dc_stat['avg_write_time_ms']:.1f}Β±{dc_stat.get('write_stddev_ms', 0):.1f}ms") - - bandwidth_summary = f" | Corrected Bandwidth: R{read_bw:.1f}Β±{dc_stat.get('read_bandwidth_stddev_mib_s', 0):.1f} W{write_bw:.1f}Β±{dc_stat.get('write_bandwidth_stddev_mib_s', 0):.1f} C{combined_bw:.1f}Β±{dc_stat.get('combined_bandwidth_stddev_mib_s', 0):.1f} MiB/s" if (read_bw > 0 or write_bw > 0 or combined_bw > 0) else "" - - self.console.print(f" β€’ {dc_name}: {completed_dc_jobs}/{total_dc_jobs} jobs " - f"({failed_dc_jobs} failed, {dc_success_rate:.1f}% success rate){timing_summary}{bandwidth_summary}") - - # Enhanced worker performance completion summary with corrected bandwidth - if worker_stats: - self.console.print(f"[bright_blue]πŸ‘₯ Enhanced Worker Performance Summary with Corrected Bandwidth:[/bright_blue]") - total_worker_records = sum(w['records_written'] for w in worker_stats.values()) - avg_worker_wps = sum(w['writes_per_second'] for w in worker_stats.values()) / len(worker_stats) - active_workers = len([w for w in worker_stats.values() if w['activity_level'] != 'Idle']) - total_read_bw = sum(w.get('avg_read_bandwidth_mib_s', 0) for w in worker_stats.values()) - total_write_bw = sum(w.get('avg_write_bandwidth_mib_s', 0) for w in worker_stats.values()) - total_combined_bw = sum(w.get('avg_combined_bandwidth_mib_s', 0) for w in worker_stats.values()) - - # Performance distribution with corrected bandwidth - rates = [w['writes_per_second'] for w in worker_stats.values()] - read_bws = [w.get('avg_read_bandwidth_mib_s', 0) for w in worker_stats.values()] - write_bws = [w.get('avg_write_bandwidth_mib_s', 0) for w in worker_stats.values()] - combined_bws = [w.get('avg_combined_bandwidth_mib_s', 0) for w in worker_stats.values()] - - if rates and read_bws and write_bws and combined_bws: - max_rate = max(rates) - min_rate = min(rates) - rate_variance = max_rate - min_rate - max_read_bw = max(read_bws) - min_read_bw = min(read_bws) - max_write_bw = max(write_bws) - min_write_bw = min(write_bws) - max_combined_bw = max(combined_bws) - min_combined_bw = min(combined_bws) - - self.console.print(f" β€’ Total Workers: {len(worker_stats)} ({active_workers} active)") - self.console.print(f" β€’ Total Records: {total_worker_records:,}") - self.console.print(f" β€’ Average Rate: {avg_worker_wps:.1f} writes/s per worker") - self.console.print(f" β€’ System Corrected Bandwidth: Read {total_read_bw:.1f} MiB/s, Write {total_write_bw:.1f} MiB/s, Combined {total_combined_bw:.1f} MiB/s") - if rates and read_bws and write_bws and combined_bws: - self.console.print(f" β€’ Performance Range: {min_rate:.1f} - {max_rate:.1f} writes/s (variance: {rate_variance:.1f})") - self.console.print(f" β€’ Corrected Bandwidth Range: Read {min_read_bw:.1f} - {max_read_bw:.1f} MiB/s, Write {min_write_bw:.1f} - {max_write_bw:.1f} MiB/s, Combined {min_combined_bw:.1f} - {max_combined_bw:.1f} MiB/s") - - # Show Enhanced ZMQ specific statistics for this group with corrected bandwidth - tasks_sent_to_zmq = group_stats.get('tasks_sent_to_zmq', 0) - records_aggregated = group_stats.get('records_aggregated', 0) - files_grouped = group_stats.get('files_grouped', 0) - avg_tasks_per_job = group_stats.get('avg_tasks_per_job', 0) - - self.console.print(f"[bright_blue]πŸ”Œ Enhanced ZMQ Performance with Corrected Bandwidth: {tasks_sent_to_zmq} tasks sent, " - f"{records_aggregated} records aggregated, " - f"{files_grouped} files grouped " - f"({avg_tasks_per_job:.1f} avg tasks/job) | " - f"Total System Bandwidth: {combined_system_bw:.1f} MiB/s[/bright_blue]") - - # Show ZMQ error summary - zmq_send_errors = group_stats.get('zmq_send_errors', 0) - zmq_recv_errors = group_stats.get('zmq_recv_errors', 0) - if zmq_send_errors == 0 and zmq_recv_errors == 0: - self.console.print("[green]🎯 Perfect ZMQ message delivery: Zero errors![/green]") - else: - self.console.print(f"[yellow]⚠️ ZMQ Message Errors: {zmq_send_errors} send, {zmq_recv_errors} recv[/yellow]") - - # ENHANCED FINAL SUMMARY WITH ZMQ, PER-DATACENTER, PER-WORKER METRICS, AND CORRECTED BANDWIDTH - total_time = time.time() - start_time - - self.console.print(f"\n[bold green]πŸŽ‰ Enhanced ZMQ PUSH/PULL WARC processing with per-datacenter, worker performance, and corrected bandwidth monitoring completed![/bold green]") - if resume: - self.console.print(f"[cyan]πŸ”„ Offset-level resume mode successfully completed[/cyan]") - - # Enhanced ZMQ processing summary with corrected bandwidth - if all_final_stats: - self.console.print(f"\nπŸ”Œ [bold]Enhanced ZMQ PUSH/PULL Architecture Summary with Corrected Bandwidth:[/bold]") - self.console.print(f" β€’ Transport: inproc:// (in-process)") - self.console.print(f" β€’ Sockets: PUSH/PULL for reliable delivery, PUB/SUB for telemetry") - self.console.print(f" β€’ High-water mark: {zmq_hwm} (automatic back-pressure)") - self.console.print(f" β€’ Records sent to ZMQ: {all_final_stats.get('tasks_sent_to_zmq', 0):,}") - self.console.print(f" β€’ Records aggregated: {all_final_stats.get('records_aggregated', 0):,}") - self.console.print(f" β€’ Jobs created: {all_final_stats.get('jobs_created', 0):,}") - self.console.print(f" β€’ Files grouped: {all_final_stats.get('files_grouped', 0):,}") - self.console.print(f" β€’ Average tasks per job: {all_final_stats.get('avg_tasks_per_job', 0):.1f}") - combined_system_bandwidth = all_final_stats.get('combined_system_bandwidth_mib_s', 0) - combined_read_bandwidth = all_final_stats.get('combined_read_bandwidth_mib_s', 0) - combined_write_bandwidth = all_final_stats.get('combined_write_bandwidth_mib_s', 0) - self.console.print(f" β€’ System Corrected Bandwidth: Read {combined_read_bandwidth:.1f} MiB/s, Write {combined_write_bandwidth:.1f} MiB/s, Combined {combined_system_bandwidth:.1f} MiB/s") - - # Enhanced per-datacenter final summary with corrected bandwidth and expanded timing - dc_stats = all_final_stats.get("datacenter_stats", {}) - if dc_stats: - self.console.print(f"\n🏒 [bold]Final Per-Datacenter Summary with Enhanced Performance Analysis & Corrected Bandwidth:[/bold]") - total_dc_jobs = sum(dc["total_jobs"] for dc in dc_stats.values()) - total_dc_completed = sum(dc["jobs_completed"] for dc in dc_stats.values()) - total_dc_failed = sum(dc["jobs_failed"] for dc in dc_stats.values()) - - for dc_name, dc_stat in dc_stats.items(): - dc_pct_of_total = (dc_stat["total_jobs"] / max(total_dc_jobs, 1)) * 100 - - # Enhanced performance metrics with corrected bandwidth - timing_details = "" - if "avg_seek_time_ms" in dc_stat and dc_stat["avg_seek_time_ms"] > 0: - timing_details = (f" | Performance: O{dc_stat.get('avg_file_open_time_ms', 0):.1f}Β±{dc_stat.get('file_open_stddev_ms', 0):.1f} " - f"S{dc_stat['avg_seek_time_ms']:.1f}Β±{dc_stat.get('seek_stddev_ms', 0):.1f} " - f"R{dc_stat['avg_read_time_ms']:.1f}Β±{dc_stat.get('read_stddev_ms', 0):.1f} " - f"W{dc_stat['avg_write_time_ms']:.1f}Β±{dc_stat.get('write_stddev_ms', 0):.1f}ms") - - bandwidth_details = "" - read_bw = dc_stat.get("avg_read_bandwidth_mib_s", 0) - write_bw = dc_stat.get("avg_write_bandwidth_mib_s", 0) - combined_bw = dc_stat.get("avg_combined_bandwidth_mib_s", 0) - if read_bw > 0 or write_bw > 0 or combined_bw > 0: - bandwidth_details = f" | Corrected Bandwidth: R{read_bw:.1f}Β±{dc_stat.get('read_bandwidth_stddev_mib_s', 0):.1f} W{write_bw:.1f}Β±{dc_stat.get('write_bandwidth_stddev_mib_s', 0):.1f} C{combined_bw:.1f}Β±{dc_stat.get('combined_bandwidth_stddev_mib_s', 0):.1f} MiB/s" - - self.console.print(f" β€’ {dc_name}: {dc_stat['jobs_completed']}/{dc_stat['total_jobs']} jobs " - f"({dc_stat['jobs_failed']} failed, {dc_stat['success_rate']:.1f}% success, " - f"{dc_pct_of_total:.1f}% of workload){timing_details}{bandwidth_details}") - - overall_dc_success = (total_dc_completed / max(total_dc_completed + total_dc_failed, 1)) * 100 - self.console.print(f" β€’ Overall: {total_dc_completed}/{total_dc_jobs} jobs " - f"({total_dc_failed} failed, {overall_dc_success:.1f}% success rate)") - - # Enhanced per-worker final summary with detailed performance analysis and corrected bandwidth - worker_stats = all_final_stats.get("worker_stats", {}) - if worker_stats: - self.console.print(f"\nπŸ‘₯ [bold]Final Worker Performance Analysis with Corrected Bandwidth:[/bold]") - - total_worker_records = sum(w['records_written'] for w in worker_stats.values()) - worker_rates = [w['writes_per_second'] for w in worker_stats.values()] - worker_read_bws = [w.get('avg_read_bandwidth_mib_s', 0) for w in worker_stats.values()] - worker_write_bws = [w.get('avg_write_bandwidth_mib_s', 0) for w in worker_stats.values()] - worker_combined_bws = [w.get('avg_combined_bandwidth_mib_s', 0) for w in worker_stats.values()] - active_workers = len([w for w in worker_stats.values() if w['activity_level'] != 'Idle']) - - if worker_rates and worker_read_bws and worker_write_bws and worker_combined_bws: - avg_rate = statistics.mean(worker_rates) - max_rate = max(worker_rates) - min_rate = min(worker_rates) - rate_stddev = statistics.stdev(worker_rates) if len(worker_rates) > 1 else 0.0 - - avg_read_bw = statistics.mean(worker_read_bws) - avg_write_bw = statistics.mean(worker_write_bws) - avg_combined_bw = statistics.mean(worker_combined_bws) - max_read_bw = max(worker_read_bws) - min_read_bw = min(worker_read_bws) - max_write_bw = max(worker_write_bws) - min_write_bw = min(worker_write_bws) - max_combined_bw = max(worker_combined_bws) - min_combined_bw = min(worker_combined_bws) - total_system_combined_bw = sum(worker_combined_bws) - - self.console.print(f" β€’ Workers: {len(worker_stats)} total ({active_workers} active)") - self.console.print(f" β€’ Total Records: {total_worker_records:,}") - self.console.print(f" β€’ Average Rate: {avg_rate:.1f}Β±{rate_stddev:.1f} writes/s per worker") - self.console.print(f" β€’ Performance Range: {min_rate:.1f} - {max_rate:.1f} writes/s") - self.console.print(f" β€’ Average Corrected Bandwidth: Read {avg_read_bw:.1f} MiB/s, Write {avg_write_bw:.1f} MiB/s, Combined {avg_combined_bw:.1f} MiB/s per worker") - self.console.print(f" β€’ Corrected Bandwidth Range: Read {min_read_bw:.1f} - {max_read_bw:.1f} MiB/s, Write {min_write_bw:.1f} - {max_write_bw:.1f} MiB/s, Combined {min_combined_bw:.1f} - {max_combined_bw:.1f} MiB/s") - self.console.print(f" β€’ Total System Corrected Bandwidth: {total_system_combined_bw:.1f} MiB/s") - - # Efficiency analysis with corrected bandwidth - theoretical_max = max_rate * len(worker_stats) - actual_combined = sum(worker_rates) - efficiency = (actual_combined / max(theoretical_max, 1)) * 100 - self.console.print(f" β€’ Worker Efficiency: {efficiency:.1f}% of theoretical maximum") - - # Show top and bottom performers with corrected bandwidth - sorted_workers = sorted(worker_stats.items(), key=lambda x: x[1]['writes_per_second'], reverse=True) - top_3 = sorted_workers[:3] if len(sorted_workers) >= 3 else sorted_workers - bottom_3 = sorted_workers[-3:] if len(sorted_workers) >= 3 else [] - - if top_3: - top_performers = ", ".join(f"T{w[0]}({w[1]['writes_per_second']:.1f}/s, C{w[1].get('avg_combined_bandwidth_mib_s', 0):.1f}MiB/s)" for w in top_3) - self.console.print(f" β€’ Top Performers: {top_performers}") - - if bottom_3 and len(sorted_workers) > 3: - bottom_performers = ", ".join(f"T{w[0]}({w[1]['writes_per_second']:.1f}/s, C{w[1].get('avg_combined_bandwidth_mib_s', 0):.1f}MiB/s)" for w in bottom_3) - self.console.print(f" β€’ Need Optimization: {bottom_performers}") - - # ZMQ reliability assessment with corrected bandwidth considerations - zmq_send_errors = all_final_stats.get('zmq_send_errors', 0) - zmq_recv_errors = all_final_stats.get('zmq_recv_errors', 0) - total_zmq_messages = all_final_stats.get('tasks_sent_to_zmq', 0) + all_final_stats.get('jobs_created', 0) - - if zmq_send_errors == 0 and zmq_recv_errors == 0: - self.console.print(f" β€’ [green]🌟 PERFECT MESSAGE DELIVERY: Zero errors across {total_zmq_messages:,} messages[/green]") - else: - error_rate = ((zmq_send_errors + zmq_recv_errors) / max(total_zmq_messages, 1)) * 100 - self.console.print(f" β€’ [yellow]πŸ“Š Message Error Rate: {error_rate:.3f}% ({zmq_send_errors + zmq_recv_errors} errors)[/yellow]") - - # Thread health summary - aggregator_alive = all_final_stats.get('aggregator_alive', False) - executor_alive = all_final_stats.get('executor_alive', False) - stats_collector_alive = all_final_stats.get('stats_collector_alive', False) - - thread_health = [aggregator_alive, executor_alive, stats_collector_alive] - healthy_threads = sum(thread_health) - - self.console.print(f" β€’ Thread Health: {healthy_threads}/3 threads healthy " - f"(Aggregator: {'βœ…' if aggregator_alive else '❌'}, " - f"Executor: {'βœ…' if executor_alive else '❌'}, " - f"StatsCollector: {'βœ…' if stats_collector_alive else '❌'})") - - # Enhanced architecture benefits achieved with corrected bandwidth monitoring - self.console.print(f"\nπŸ—οΈ [bold]Enhanced ZMQ Architecture Benefits Achieved:[/bold]") - self.console.print(f" β€’ βœ… Automatic back-pressure: HWM prevents memory overflow") - self.console.print(f" β€’ βœ… Reliable message delivery: PUSH/PULL guarantees no message loss") - self.console.print(f" β€’ βœ… Non-blocking telemetry: PUB/SUB prevents stats from blocking pipeline") - self.console.print(f" β€’ βœ… Decoupled processing: Independent producer/consumer scaling") - self.console.print(f" β€’ βœ… Future-proof scaling: Easy upgrade to ipc:// or tcp:// transport") - self.console.print(f" β€’ βœ… Zero task loss: Mathematical guarantee maintained") - self.console.print(f" β€’ βœ… Per-datacenter progress: Detailed monitoring and ETA calculation") - self.console.print(f" β€’ βœ… Per-worker performance: Individual thread monitoring with timing analysis") - self.console.print(f" β€’ βœ… Comprehensive timing: Seek/read/write performance with standard deviations") - self.console.print(f" β€’ βœ… Queue growth adaptation: Smart ETA accounting for dynamic workloads") - self.console.print(f" β€’ βœ… CORRECTED: Real-time bandwidth monitoring: Proper MiB/s tracking with variance analysis") - self.console.print(f" β€’ βœ… CORRECTED: Expanded timing display: Four separate columns (Open/Seek/Read/Write Β±Οƒ)") - self.console.print(f" β€’ βœ… CORRECTED: Performance consistency analysis: Proper statistical variance tracking across workers and datacenters") - - if resume: - self.console.print(f" β€’ βœ… Offset-level resume: Perfect granularity preserved") - - # Final success assessment with enhanced criteria including corrected bandwidth - task_loss = all_final_stats.get('task_loss', 0) - success_rate = all_final_stats.get('task_success_rate', 0) - zmq_message_reliability = (zmq_send_errors + zmq_recv_errors) == 0 - combined_bandwidth = all_final_stats.get('combined_system_bandwidth_mib_s', 0) - - # Worker performance assessment with corrected bandwidth - worker_performance_good = True - bandwidth_performance_good = True - if worker_stats: - worker_rates = [w['writes_per_second'] for w in worker_stats.values()] - worker_combined_bandwidths = [w.get('avg_combined_bandwidth_mib_s', 0) for w in worker_stats.values()] - - if worker_rates: - rate_stddev = statistics.stdev(worker_rates) if len(worker_rates) > 1 else 0.0 - avg_rate = statistics.mean(worker_rates) - coefficient_of_variation = (rate_stddev / max(avg_rate, 1)) * 100 - worker_performance_good = coefficient_of_variation < 50 # Less than 50% variation is good + # Note: Detailed statistics will be printed via the processor's stored report - if worker_combined_bandwidths: - bw_stddev = statistics.stdev(worker_combined_bandwidths) if len(worker_combined_bandwidths) > 1 else 0.0 - avg_bw = statistics.mean(worker_combined_bandwidths) - bw_coefficient_of_variation = (bw_stddev / max(avg_bw, 1)) * 100 if avg_bw > 0 else 0 - bandwidth_performance_good = bw_coefficient_of_variation < 50 # Less than 50% variation is good - - if (task_loss == 0 and success_rate > 95 and zmq_message_reliability and - worker_performance_good and bandwidth_performance_good and combined_bandwidth > 0): - self.console.print(f"\n[green]🌟 OUTSTANDING: Perfect processing with {success_rate:.1f}% success rate, " - f"zero task loss, perfect ZMQ message delivery, excellent worker & bandwidth consistency " - f"({combined_bandwidth:.1f} MiB/s), comprehensive per-datacenter tracking with expanded timing columns![/green]") - elif task_loss == 0 and success_rate > 95 and zmq_message_reliability and combined_bandwidth > 0: - self.console.print(f"\n[green]βœ… EXCELLENT: Perfect processing with {success_rate:.1f}% success rate, " - f"perfect ZMQ message delivery, corrected bandwidth monitoring ({combined_bandwidth:.1f} MiB/s), " - f"enhanced per-datacenter & worker architecture with expanded timing![/green]") - elif task_loss == 0 and success_rate > 95: - self.console.print(f"\n[green]βœ… VERY GOOD: Perfect processing with {success_rate:.1f}% success rate " - f"and comprehensive ZMQ per-datacenter, worker & corrected bandwidth monitoring![/green]") - elif task_loss == 0: - self.console.print(f"\n[green]βœ… GOOD: Zero task loss with Enhanced ZMQ PUSH/PULL per-datacenter, worker & corrected bandwidth architecture[/green]") - - # Processing rate comparison with worker efficiency and corrected bandwidth - if total_time > 0: - records_per_second = all_final_stats.get('records_written', 0) / total_time - self.console.print(f"\nπŸ“ˆ [bold]Enhanced Performance Summary with Corrected Bandwidth:[/bold]") - self.console.print(f" β€’ Total runtime: {total_time:.1f} seconds") - self.console.print(f" β€’ Processing rate: {records_per_second:.1f} records/second") - self.console.print(f" β€’ File grouping efficiency: {all_final_stats.get('avg_tasks_per_job', 0):.1f} tasks per job") - self.console.print(f" β€’ System corrected bandwidth: {combined_bandwidth:.1f} MiB/s combined") - - if worker_stats: - total_worker_records = sum(w['records_written'] for w in worker_stats.values()) - combined_worker_rate = sum(w['writes_per_second'] for w in worker_stats.values()) - combined_worker_combined_bw = sum(w.get('avg_combined_bandwidth_mib_s', 0) for w in worker_stats.values()) - self.console.print(f" β€’ Worker throughput: {combined_worker_rate:.1f} combined writes/s from {len(worker_stats)} workers") - self.console.print(f" β€’ Worker corrected bandwidth: {combined_worker_combined_bw:.1f} combined MiB/s from {len(worker_stats)} workers") - self.console.print(f" β€’ Parallelization factor: {combined_worker_rate / max(records_per_second, 1):.1f}x") + # Print the comprehensive report from the processor + if processor: + try: + self.console.print("\n" + "="*80) + processor.print_completion_report() + self.console.print("="*80) + except Exception as e: + self.console.print(f"[red]❌ Could not display completion report: {e}[/red]") + # Fallback to basic summary + if all_final_stats: + success_rate = all_final_stats.get('task_success_rate', 0) + jobs_completed = all_final_stats.get('jobs_completed', 0) + records_written = all_final_stats.get('records_written', 0) + combined_bandwidth = all_final_stats.get('combined_system_bandwidth_mib_s', 0) + + self.console.print(f"[green]πŸŽ‰ Processing completed successfully![/green]") + self.console.print(f"[green]πŸ“Š Final Summary: {jobs_completed} jobs, {records_written} records, " + f"{success_rate:.1f}% success rate, {combined_bandwidth:.1f} MiB/s bandwidth[/green]") return CommandResult( success=True, object={ **all_final_stats, # Include all Enhanced ZMQ stats with corrected bandwidth - "architecture": "enhanced_zmq_push_pull_with_pub_sub_telemetry_per_datacenter_worker_performance_and_corrected_bandwidth_monitoring", - "zmq_architecture": True, - "corrected_bandwidth_monitoring": True, # CORRECTED - "expanded_timing_columns": True, - "performance_variance_analysis": True, - "transport": "inproc", - "reliable_message_delivery": zmq_message_reliability, - "automatic_back_pressure": True, - "decoupled_processing": True, - "per_datacenter_progress": True, - "per_worker_performance": True, - "comprehensive_timing_analysis": True, - "corrected_real_time_bandwidth_analysis": True, # CORRECTED - "statistical_variance_tracking": True, "configuration": { "max_workers": max_workers, "record_threshold": record_threshold, @@ -3125,52 +2647,18 @@ def warc(self, local_specifier: str, remote_specifier: str, urls_file: str = Non "rollover_limit": rollover_limit, "resume_mode": resume, "stats_interval": stats_interval - }, - "zmq_statistics": { - "total_messages_sent": all_final_stats.get('tasks_sent_to_zmq', 0) + all_final_stats.get('jobs_created', 0), - "send_errors": zmq_send_errors, - "recv_errors": zmq_recv_errors, - "message_reliability_percent": ((total_zmq_messages - zmq_send_errors - zmq_recv_errors) / max(total_zmq_messages, 1)) * 100, - "thread_health": { - "aggregator_alive": aggregator_alive, - "executor_alive": executor_alive, - "stats_collector_alive": stats_collector_alive, - "healthy_thread_count": healthy_threads - } - }, - "datacenter_statistics": all_final_stats.get("datacenter_stats", {}), - "worker_performance_statistics": all_final_stats.get("worker_stats", {}), - # CORRECTED: Enhanced performance analysis with proper bandwidth - "enhanced_performance_analysis": { - "worker_performance_consistent": worker_performance_good, - "bandwidth_performance_consistent": bandwidth_performance_good, # CORRECTED - "total_active_workers": len([w for w in worker_stats.values() if w['activity_level'] != 'Idle']) if worker_stats else 0, - "combined_system_bandwidth_mib_s": combined_bandwidth, # CORRECTED - "worker_efficiency_analysis": { - "average_rate": statistics.mean([w['writes_per_second'] for w in worker_stats.values()]) if worker_stats else 0, - "rate_standard_deviation": statistics.stdev([w['writes_per_second'] for w in worker_stats.values()]) if worker_stats and len(worker_stats) > 1 else 0, - "coefficient_of_variation": (statistics.stdev([w['writes_per_second'] for w in worker_stats.values()]) / max(statistics.mean([w['writes_per_second'] for w in worker_stats.values()]), 1)) * 100 if worker_stats and len(worker_stats) > 1 else 0, - # CORRECTED: Bandwidth efficiency analysis using combined bandwidth - "average_combined_bandwidth_mib_s": statistics.mean([w.get('avg_combined_bandwidth_mib_s', 0) for w in worker_stats.values()]) if worker_stats else 0, - "combined_bandwidth_standard_deviation": statistics.stdev([w.get('avg_combined_bandwidth_mib_s', 0) for w in worker_stats.values()]) if worker_stats and len(worker_stats) > 1 else 0, - "combined_bandwidth_coefficient_of_variation": (statistics.stdev([w.get('avg_combined_bandwidth_mib_s', 0) for w in worker_stats.values()]) / max(statistics.mean([w.get('avg_combined_bandwidth_mib_s', 0) for w in worker_stats.values()]), 1)) * 100 if worker_stats and len(worker_stats) > 1 and statistics.mean([w.get('avg_combined_bandwidth_mib_s', 0) for w in worker_stats.values()]) > 0 else 0 - } if worker_stats else {} } }, - msg=f"Successfully processed {total_processed} jobs with Enhanced ZMQ PUSH/PULL Per-Datacenter, Worker Performance & Corrected Bandwidth Monitoring Architecture " - f"(Success rate: {all_final_stats.get('task_success_rate', 0):.1f}%, " - f"Message reliability: {((total_zmq_messages - zmq_send_errors - zmq_recv_errors) / max(total_zmq_messages, 1)) * 100:.2f}%, " - f"Records sent to ZMQ: {all_final_stats.get('tasks_sent_to_zmq', 0):,}, " - f"Files grouped: {all_final_stats.get('files_grouped', 0):,}, " - f"Datacenters monitored: {len(all_final_stats.get('datacenter_stats', {}))}, " - f"Workers monitored: {len(all_final_stats.get('worker_stats', {}))}, " - f"Corrected system bandwidth: {combined_bandwidth:.1f} MiB/s, " - f"Zero task loss: {all_final_stats.get('task_loss', 0) == 0}, " - f"Thread health: {healthy_threads}/3, " - f"Worker performance: {'Consistent' if worker_performance_good else 'Variable'}, " - f"Bandwidth consistency: {'Consistent' if bandwidth_performance_good else 'Variable'})" + msg=f"Successfully processed {total_processed} jobs with Enhanced ZMQ PUSH/PULL Per-Datacenter, Worker Performance & Corrected Bandwidth Monitoring Architecture" ) except Exception as e: + # Print error report if processor is available + if 'processor' in locals() and processor: + try: + processor.print_completion_report() + except: + pass + self.console.print(f"[red]❌ Enhanced ZMQ PUSH/PULL per-datacenter, worker & corrected bandwidth processing failed: {e}[/red]") raise \ No newline at end of file diff --git a/query_warc.py b/query_warc.py new file mode 100644 index 0000000..1ef8423 --- /dev/null +++ b/query_warc.py @@ -0,0 +1,35 @@ + for dc_name, dc_stat in dc_stats.items(): + dc_pct_of_total = (dc_stat["total_jobs"] / max(total_dc_jobs, 1)) * 100 + + # Performance metrics + timing_summary = f"{dc_stat.get('avg_file_open_time_ms', 0):.1f}/{dc_stat.get('avg_seek_time_ms', 0):.1f}/{dc_stat.get('avg_read_time_ms', 0):.1f}/{dc_stat.get('avg_write_time_ms', 0):.1f}" + + # Bandwidth metrics + read_bw = dc_stat.get("avg_read_bandwidth_mib_s", 0) + write_bw = dc_stat.get("avg_write_bandwidth_mib_s", 0) + combined_bw = dc_stat.get("total_bytes_mb", 0) / dc_stat.get("elapsed", 0) if dc_stat.get("elapsed", 0) > 0 else 0 + bandwidth_summary = f"{read_bw:.1f}/{write_bw:.1f}/{combined_bw:.1f}" + + job_summary = f"{dc_stat['jobs_completed']}/{dc_stat['total_jobs']} ({dc_stat['jobs_failed']} failed)" + + md.append(f"| {dc_name} | {job_summary} | {dc_stat['success_rate']:.1f}% | {timing_summary} | {bandwidth_summary} |") + + # Calculate aggregated performance and bandwidth metrics + total_open_time = sum(dc.get('avg_file_open_time_ms', 0) * dc.get('total_jobs', 0) for dc in dc_stats.values()) / max(total_dc_jobs, 1) + total_seek_time = sum(dc.get('avg_seek_time_ms', 0) * dc.get('total_jobs', 0) for dc in dc_stats.values()) / max(total_dc_jobs, 1) + total_read_time = sum(dc.get('avg_read_time_ms', 0) * dc.get('total_jobs', 0) for dc in dc_stats.values()) / max(total_dc_jobs, 1) + total_write_time = sum(dc.get('avg_write_time_ms', 0) * dc.get('total_jobs', 0) for dc in dc_stats.values()) / max(total_dc_jobs, 1) + + total_read_bw = sum(dc.get('avg_read_bandwidth_mib_s', 0) * dc.get('total_jobs', 0) for dc in dc_stats.values()) / max(total_dc_jobs, 1) + total_write_bw = sum(dc.get('avg_write_bandwidth_mib_s', 0) * dc.get('total_jobs', 0) for dc in dc_stats.values()) / max(total_dc_jobs, 1) + + total_bytes = sum(dc.get('total_bytes_mb', 0) for dc in dc_stats.values()) + total_elapsed = sum(dc.get('elapsed', 0) * dc.get('total_jobs', 0) for dc in dc_stats.values()) / max(total_dc_jobs, 1) + total_combined_bw = total_bytes / max(total_elapsed, 0.001) + + total_timing_summary = f"{total_open_time:.1f}/{total_seek_time:.1f}/{total_read_time:.1f}/{total_write_time:.1f}" + total_bandwidth_summary = f"{total_read_bw:.1f}/{total_write_bw:.1f}/{total_combined_bw:.1f}" + + overall_dc_success = (total_dc_completed / max(total_dc_completed + total_dc_failed, 1)) * 100 + md.append(f"| **TOTAL** | **{total_dc_completed}/{total_dc_jobs}** | **{overall_dc_success:.1f}%** | **{total_timing_summary}** | **{total_bandwidth_summary}** |") + md.append("")