NikhilPatil/Predictive_Maintanance_With_RAG_And_Fleet_Simulation
0
1import numpy as np2import pandas as pd3 4 5def get_sensor_cols(df):6 return [f"s{i}" for i in range(1, 22) if f"s{i}" in df.columns]7 8 9def get_op_cols(df):10 return [c for c in ["op1", "op2", "op3"] if c in df.columns]11 12 13def safe_numeric_columns(df, exclude=None):14 if exclude is None:15 exclude = []16 17 cols = []18 19 for c in df.columns:20 if c in exclude:21 continue22 23 if pd.api.types.is_numeric_dtype(df[c]):24 cols.append(c)25 26 return cols27 28 29def detect_constant_like_columns(df, cols, std_threshold=1e-8, unique_threshold=1):30 constant_cols = []31 32 for c in cols:33 if c not in df.columns:34 continue35 36 std_val = float(df[c].std()) if df[c].notna().any() else 0.037 nunique = int(df[c].nunique(dropna=True))38 39 if std_val <= std_threshold or nunique <= unique_threshold:40 constant_cols.append(c)41 42 return constant_cols43 44 45def apply_small_sensor_noise(46 block,47 reference_df,48 noise_cols,49 constant_like_cols,50 rng,51 noise_scale=0.00352):53 """54 Adds very small sensor noise without touching unit/cycle/target.55 Noise is intentionally tiny because this is inference/testing data,56 not training augmentation.57 """58 59 block = block.copy()60 61 usable_cols = [62 c for c in noise_cols63 if c in block.columns and c not in constant_like_cols64 ]65 66 for c in usable_cols:67 ref = reference_df[c].dropna().astype(float)68 69 if len(ref) < 5:70 continue71 72 std_val = float(ref.std())73 74 if std_val <= 0:75 continue76 77 noise = rng.normal(78 loc=0.0,79 scale=noise_scale * std_val,80 size=len(block)81 )82 83 block[c] = block[c].astype(float) + noise84 85 original_non_na = reference_df[c].dropna()86 87 if len(original_non_na) > 0:88 is_integer_like = np.all(89 np.isclose(original_non_na, np.round(original_non_na))90 )91 92 if is_integer_like:93 block[c] = np.round(block[c]).astype(int)94 95 return block96 97 98def validate_sequential_sample(original_df, sampled_df, target_col="target"):99 sensor_cols = get_sensor_cols(original_df)100 op_cols = get_op_cols(original_df)101 102 compare_cols = [103 c104 for c in op_cols + sensor_cols + ([target_col] if target_col in original_df.columns else [])105 if c in original_df.columns and c in sampled_df.columns106 ]107 108 rows = []109 110 for c in compare_cols:111 orig = original_df[c].dropna().astype(float)112 samp = sampled_df[c].dropna().astype(float)113 114 if len(orig) == 0 or len(samp) == 0:115 continue116 117 rows.append({118 "feature": c,119 "original_mean": float(orig.mean()),120 "sampled_mean": float(samp.mean()),121 "mean_diff": float(samp.mean() - orig.mean()),122 "original_std": float(orig.std()),123 "sampled_std": float(samp.std()),124 "std_diff": float(samp.std() - orig.std()),125 "original_median": float(orig.median()),126 "sampled_median": float(samp.median()),127 "median_diff": float(samp.median() - orig.median()),128 "original_min": float(orig.min()),129 "sampled_min": float(samp.min()),130 "original_max": float(orig.max()),131 "sampled_max": float(samp.max())132 })133 134 validation_df = pd.DataFrame(rows)135 136 original_life = original_df.groupby("unit")["cycle"].max()137 sampled_life = sampled_df.groupby("unit")["cycle"].max()138 139 summary = {140 "original_rows": int(len(original_df)),141 "sampled_rows": int(len(sampled_df)),142 "original_units": int(original_df["unit"].nunique()),143 "sampled_units": int(sampled_df["unit"].nunique()),144 "original_lifespan_mean": float(original_life.mean()),145 "sampled_lifespan_mean": float(sampled_life.mean()),146 "original_lifespan_median": float(original_life.median()),147 "sampled_lifespan_median": float(sampled_life.median()),148 "original_lifespan_min": int(original_life.min()),149 "sampled_lifespan_min": int(sampled_life.min()),150 "original_lifespan_max": int(original_life.max()),151 "sampled_lifespan_max": int(sampled_life.max())152 }153 154 return summary, validation_df155 156 157def create_distribution_preserving_sequential_sample(158 input_df,159 num_sampled_engines=None,160 random_seed=42,161 target_col="target",162 sampling_strategy="full_engine_or_prefix",163 min_block_length=30,164 max_block_length=None,165 add_noise=False,166 noise_scale=0.003,167 new_unit_start=10001168):169 """170 Corrected sequential sampler.171 172 Main rule:173 Do not break sensor-target relationship.174 175 This sampler samples:176 - full engine trajectories177 - or prefixes from the beginning of engine trajectories178 179 It preserves:180 - cycle order181 - original cycle values182 - original target values183 - sensor degradation pattern184 - target-sensor alignment185 186 It does NOT:187 - randomly sample rows188 - reindex cycle to 1..N189 - recompute target as block_length - cycle190 - create artificial failure labels191 """192 193 if "unit" not in input_df.columns or "cycle" not in input_df.columns:194 raise ValueError("input_df must contain unit and cycle columns.")195 196 rng = np.random.default_rng(random_seed)197 198 original_columns = list(input_df.columns)199 200 df = input_df.copy()201 df = df.sort_values(["unit", "cycle"]).reset_index(drop=True)202 203 engine_groups = {204 int(unit): g.copy().sort_values("cycle").reset_index(drop=True)205 for unit, g in df.groupby("unit")206 }207 208 if num_sampled_engines is None:209 num_sampled_engines = len(engine_groups)210 211 numeric_noise_cols = safe_numeric_columns(212 df,213 exclude=["unit", "cycle", target_col, "target_raw"]214 )215 216 constant_like_cols = detect_constant_like_columns(217 df,218 numeric_noise_cols219 )220 221 sampled_blocks = []222 metadata_records = []223 224 unit_ids = list(engine_groups.keys())225 226 for i in range(int(num_sampled_engines)):227 new_unit_id = int(new_unit_start + i)228 229 source_unit = int(rng.choice(unit_ids))230 source_group = engine_groups[source_unit].copy()231 source_group = source_group.sort_values("cycle").reset_index(drop=True)232 233 n = len(source_group)234 235 if n < min_block_length:236 selected = source_group.copy()237 chosen_strategy = "full_engine_short"238 239 else:240 if sampling_strategy == "full_engine":241 selected = source_group.copy()242 chosen_strategy = "full_engine"243 244 elif sampling_strategy == "prefix":245 max_len = n246 247 if max_block_length is not None:248 max_len = min(max_len, int(max_block_length))249 250 prefix_len = int(rng.integers(int(min_block_length), max_len + 1))251 selected = source_group.iloc[:prefix_len].copy()252 chosen_strategy = "prefix"253 254 elif sampling_strategy == "full_engine_or_prefix":255 chosen_strategy = rng.choice(["full_engine", "prefix"], p=[0.55, 0.45])256 257 if chosen_strategy == "full_engine":258 selected = source_group.copy()259 260 else:261 max_len = n262 263 if max_block_length is not None:264 max_len = min(max_len, int(max_block_length))265 266 prefix_len = int(rng.integers(int(min_block_length), max_len + 1))267 selected = source_group.iloc[:prefix_len].copy()268 269 else:270 raise ValueError(271 "sampling_strategy must be full_engine, prefix, or full_engine_or_prefix."272 )273 274 original_start_cycle = int(selected["cycle"].min())275 original_end_cycle = int(selected["cycle"].max())276 277 selected = selected.copy()278 selected["unit"] = new_unit_id279 280 # Important:281 # Do NOT reindex cycle.282 # Do NOT recompute target.283 # Preserve original cycle and target alignment.284 285 if add_noise:286 selected = apply_small_sensor_noise(287 block=selected,288 reference_df=df,289 noise_cols=numeric_noise_cols,290 constant_like_cols=constant_like_cols,291 rng=rng,292 noise_scale=float(noise_scale)293 )294 295 sampled_blocks.append(selected)296 297 metadata_records.append({298 "new_unit": new_unit_id,299 "source_unit": source_unit,300 "sampling_strategy": str(chosen_strategy),301 "sampled_length": int(len(selected)),302 "original_start_cycle": original_start_cycle,303 "original_end_cycle": original_end_cycle,304 "new_start_cycle": int(selected["cycle"].min()),305 "new_end_cycle": int(selected["cycle"].max()),306 "target_preserved": bool(target_col in selected.columns),307 "cycle_preserved": True,308 "noise_added": bool(add_noise),309 "noise_scale": float(noise_scale)310 })311 312 sampled_df = pd.concat(sampled_blocks, ignore_index=True)313 sampled_df = sampled_df.sort_values(["unit", "cycle"]).reset_index(drop=True)314 315 ordered_cols = [c for c in original_columns if c in sampled_df.columns]316 extra_cols = [c for c in sampled_df.columns if c not in ordered_cols]317 sampled_df = sampled_df[ordered_cols + extra_cols]318 319 sampling_metadata = {320 "sampling_method": "distribution_preserving_sequential_engine_prefix_sampling",321 "sampling_strategy": sampling_strategy,322 "num_sampled_engines": int(num_sampled_engines),323 "random_seed": int(random_seed),324 "min_block_length": int(min_block_length),325 "max_block_length": None if max_block_length is None else int(max_block_length),326 "add_noise": bool(add_noise),327 "noise_scale": float(noise_scale),328 "constant_like_columns": constant_like_cols,329 "important_note": (330 "This sampler preserves original cycle and target values. "331 "It does not recompute target from sampled block length."332 ),333 "records": metadata_records334 }335 336 validation_summary, validation_df = validate_sequential_sample(337 original_df=df,338 sampled_df=sampled_df,339 target_col=target_col340 )341 342 return sampled_df, sampling_metadata, validation_summary, validation_df