File size: 3,722 Bytes
fcdda81 | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 | """
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
# Load the raw record using WFDB bindings
signal_array, metadata = wfdb.rdsamp(absolute_path)
# Ensure the signal fits our target runtime input dimensions exactly
if len(signal_array) != self.signal_length:
signal_array = resample(signal_array, self.signal_length)
# Fill any missing values with zeros to maintain tensor consistency
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.
"""
# Clinical risk markers defined within PTB-XL schema guidelines
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)
# Subsample rows to fit local memory constraints during development
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}") |