| """ |
| Production Bio-Signal Telemetry Processing Engine. |
| Parses multi-lead ECG waveforms and processes diagnostic statements into clean targets. |
| """ |
| import os |
| import ast |
| import numpy as np |
| import pandas as pd |
| import wfdb |
| from scipy.signal import resample |
|
|
|
|
| class ECGPreprocessingPipeline: |
| def __init__(self, ecg_dir, target_sampling_rate=100, signal_length=1000): |
| self.ecg_dir = ecg_dir |
| self.target_sampling_rate = target_sampling_rate |
| self.signal_length = signal_length |
|
|
| def parse_diagnostic_statements(self, statement_string): |
| """Converts raw string dictionaries into clean evaluation arrays.""" |
| try: |
| return ast.literal_eval(statement_string) |
| except (ValueError, SyntaxError): |
| return {} |
|
|
| def extract_signal_waveform(self, relative_record_path): |
| """Extracts and scales 12-lead raw data points directly from binary assets.""" |
| absolute_path = os.path.join(self.ecg_dir, relative_record_path) |
| if not os.path.exists(absolute_path + ".dat"): |
| return None |
|
|
| |
| signal_array, metadata = wfdb.rdsamp(absolute_path) |
|
|
| |
| if len(signal_array) != self.signal_length: |
| signal_array = resample(signal_array, self.signal_length) |
|
|
| |
| if np.isnan(signal_array).any(): |
| signal_array = np.nan_to_num(signal_array, nan=0.0) |
|
|
| return signal_array |
|
|
| def map_diagnostic_to_binary_target(self, dict_statements): |
| """ |
| Maps multi-label classifications into clean clinical risk scores. |
| 1: High Cardiovascular Risk, 0: Normal Baseline. |
| """ |
| |
| high_risk_diagnostic_anchors = {'AMI', 'IMI', 'ALV', 'ILV', 'LVH', 'LAO/LAE', 'RVH'} |
| for diagnostic_key in dict_statements.keys(): |
| if diagnostic_key in high_risk_diagnostic_anchors: |
| return 1 |
| return 0 |
|
|
| def process_dataset(self, limits=1000): |
| metadata_csv = os.path.join(self.ecg_dir, "ptbxl_database.csv") |
| if not os.path.exists(metadata_csv): |
| print(f"[WARNING] Database tracking sheet missing at {metadata_csv}") |
| return None, None |
|
|
| df = pd.read_csv(metadata_csv, index_col='ecg_id') |
| df['scp_codes'] = df['scp_codes'].apply(self.parse_diagnostic_statements) |
| df['Target'] = df['scp_codes'].apply(self.map_diagnostic_to_binary_target) |
|
|
| |
| df_subsample = df.head(limits) |
|
|
| signal_accumulator = [] |
| target_accumulator = [] |
|
|
| print(f"[INFO] Ingesting {len(df_subsample)} raw biosignal channels from disk...") |
| for ecg_id, row in df_subsample.iterrows(): |
| waveform = self.extract_signal_waveform(row['filename_lr']) |
| if waveform is not None: |
| signal_accumulator.append(waveform) |
| target_accumulator.append(row['Target']) |
|
|
| return np.array(signal_accumulator), np.array(target_accumulator) |
|
|
|
|
| if __name__ == "__main__": |
| project_root = os.path.abspath(os.path.join(os.path.dirname(__file__), "..")) |
| ecg_data_path = os.path.join(project_root, "data", "ecg") |
|
|
| if os.path.exists(os.path.join(ecg_data_path, "ptbxl_database.csv")): |
| processor = ECGPreprocessingPipeline(ecg_data_path) |
| signals, targets = processor.process_dataset(limits=100) |
| print(f"[SUCCESS] Signal Tensor Layout: {signals.shape}, Targets Layout: {targets.shape}") |