Repository navigation
Expand file tree
/
Copy pathdata_loader.py
More file actions
365 lines (308 loc) · 11.8 KB
/
Copy pathdata_loader.py
File metadata and controls
365 lines (308 loc) · 11.8 KB
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
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
# -*- coding: utf-8 -*-
"""
Data loader for multivariate time series stored in two CSV files (GBK-encoded).
Assumptions (as requested):
- Both CSVs have headers and use GBK encoding.
- Training CSV: features start from the 5th column (index 4) to the end.
- Test CSV: features start from the 5th column (index 4) to the penultimate column;
the last column is the label (0 = normal, 1 = anomaly).
- Min-Max normalization (fit on TRAIN features only, then applied to TEST).
- Builds sliding-window PyTorch datasets/dataloaders:
* train windows are unlabeled (unsupervised)
* test windows carry a label derived from the window (default = label of the LAST row in the window)
Usage example:
from data_loader import build_dataloaders
train_loader, test_loader, meta = build_dataloaders(
train_csv='train.csv',
test_csv='test.csv',
window_size=64,
batch_size=32,
stride=1,
label_strategy='last',
drop_last=False,
)
# meta contains: feature_names, min_, max_, scaler callable, etc.
Author: you
"""
from __future__ import annotations
import os
from dataclasses import dataclass
from typing import Optional, Tuple, Dict, Any, Iterable
import numpy as np
import pandas as pd
import torch
from torch.utils.data import Dataset, DataLoader
# -----------------------------
# Utilities
# -----------------------------
def _check_file(path: str) -> None:
if not os.path.exists(path):
raise FileNotFoundError(f"File not found: {path}")
def _to_numpy(x: Any) -> np.ndarray:
if isinstance(x, np.ndarray):
return x
return np.asarray(x)
@dataclass
class MinMaxScaler:
"""Simple per-feature Min-Max scaler.
Fit on training features only; apply to any array with the same feature count.
If a feature is constant (max == min), its scaled output becomes 0.
"""
min_: np.ndarray
max_: np.ndarray
eps: float = 1e-12
def transform(self, X: np.ndarray) -> np.ndarray:
X = _to_numpy(X)
denom = np.maximum(self.max_ - self.min_, self.eps)
Xs = (X - self.min_) / denom
# Clip for numerical safety
return np.clip(Xs, 0.0, 1.0)
def inverse_transform(self, Xs: np.ndarray) -> np.ndarray:
Xs = _to_numpy(Xs)
return Xs * (self.max_ - self.min_) + self.min_
# -----------------------------
# CSV Loading & Preprocessing
# -----------------------------
def load_csvs_gbk(train_csv: str, test_csv: str,
encoding: str = 'gbk') -> Tuple[pd.DataFrame, pd.DataFrame]:
"""Load two GBK-encoded CSVs with headers."""
_check_file(train_csv)
_check_file(test_csv)
train_df = pd.read_csv(train_csv, encoding=encoding)
test_df = pd.read_csv(test_csv, encoding=encoding)
if len(train_df.columns) < 5:
raise ValueError("Training CSV must have at least 5 columns (features start from 5th column).")
if len(test_df.columns) < 6:
raise ValueError("Test CSV must have at least 6 columns (features start from 5th, last is label).")
return train_df, test_df
def split_features_labels(train_df: pd.DataFrame, test_df: pd.DataFrame) -> Tuple[np.ndarray, np.ndarray, np.ndarray, list]:
"""Extract features/labels based on the column rules.
Returns (X_train, X_test, y_test, feature_names)
"""
# Training: features = columns[4:]
X_train = train_df.iloc[:, 4:].values.astype(np.float32)
# Test: features = columns[4:-1], label = last column
X_test = test_df.iloc[:, 4:-1].values.astype(np.float32)
y_test = test_df.iloc[:, -1].values.astype(np.int64)
# Try to derive feature names consistently from the intersection of both files
feat_names_train = list(train_df.columns[4:])
feat_names_test = list(test_df.columns[4:-1])
if len(feat_names_train) != len(feat_names_test):
# fall back to test feature names length
feature_names = feat_names_test
else:
feature_names = feat_names_train
if X_train.shape[1] != X_test.shape[1]:
raise ValueError(
f"Train/Test feature dimension mismatch: train={X_train.shape[1]}, test={X_test.shape[1]}"
)
return X_train, X_test, y_test, feature_names
def fit_minmax_on_train(X_train: np.ndarray) -> MinMaxScaler:
X_train = _to_numpy(X_train)
min_ = np.nanmin(X_train, axis=0)
max_ = np.nanmax(X_train, axis=0)
# Replace NaNs if any
min_[~np.isfinite(min_)] = 0.0
max_[~np.isfinite(max_)] = 1.0
return MinMaxScaler(min_=min_.astype(np.float32), max_=max_.astype(np.float32))
# -----------------------------
# Sliding Window Dataset
# -----------------------------
class SequenceDataset(Dataset):
"""Create sliding windows from a 2D array [N, d].
For test set, optional labels (1D) can be provided. The window's label is
derived by `label_strategy`:
- 'last': label = labels[idx + T - 1]
- 'any': label = 1 if any label in the window is 1 else 0
- 'majority': label = 1 if sum(labels in window) > T/2 else 0
"""
def __init__(
self,
X: np.ndarray,
window_size: int,
stride: int = 1,
labels: Optional[np.ndarray] = None,
label_strategy: str = 'last',
device: Optional[torch.device] = None,
) -> None:
super().__init__()
X = _to_numpy(X)
if labels is not None:
labels = _to_numpy(labels)
if len(labels.shape) != 1:
raise ValueError("labels must be a 1D array.")
if len(labels) != len(X):
raise ValueError("labels length must match X rows.")
if window_size <= 0:
raise ValueError("window_size must be > 0")
if stride <= 0:
raise ValueError("stride must be > 0")
if len(X) < window_size:
raise ValueError(f"Not enough rows ({len(X)}) to make a window of size {window_size}.")
if label_strategy not in ("last", "any", "majority"):
raise ValueError("label_strategy must be one of: 'last', 'any', 'majority'")
self.X = X.astype(np.float32)
self.labels = labels.astype(np.int64) if labels is not None else None
self.T = window_size
self.stride = stride
self.label_strategy = label_strategy
self.device = device
# Pre-compute valid start indices
self.starts = list(range(0, len(self.X) - self.T + 1, self.stride))
def __len__(self) -> int:
return len(self.starts)
def _derive_label(self, start: int) -> int:
if self.labels is None:
return -1 # unused
s, e = start, start + self.T
win = self.labels[s:e]
if self.label_strategy == 'last':
return int(win[-1])
elif self.label_strategy == 'any':
return int(np.any(win == 1))
else: # 'majority'
return int(np.sum(win == 1) > (len(win) / 2.0))
def __getitem__(self, idx: int):
start = self.starts[idx]
end = start + self.T
x_win = self.X[start:end] # [T, d]
if self.labels is None:
# Train set: return only x
x_t = torch.from_numpy(x_win)
return x_t
else:
y = self._derive_label(start)
x_t = torch.from_numpy(x_win)
y_t = torch.tensor(y, dtype=torch.long)
return x_t, y_t
# -----------------------------
# Builder
# -----------------------------
def build_dataloaders(
train_csv: str,
test_csv: str,
window_size: int,
batch_size: int,
stride: int = 1,
label_strategy: str = 'last',
encoding: str = 'gbk',
num_workers: int = 0,
pin_memory: bool = True,
drop_last: bool = False,
) -> Tuple[DataLoader, DataLoader, Dict[str, Any]]:
"""Build train/test dataloaders from two CSV files.
Returns: (train_loader, test_loader, meta)
- train_loader: yields tensors [B, T, d]
- test_loader: yields (x, y) where x: [B, T, d], y: [B]
- meta: dict with feature_names, scaler, min_, max_, train_rows, test_rows
"""
train_df, test_df = load_csvs_gbk(train_csv, test_csv, encoding=encoding)
X_train, X_test, y_test, feature_names = split_features_labels(train_df, test_df)
# Fit Min-Max on train features only; apply to both
scaler = fit_minmax_on_train(X_train)
X_train_s = scaler.transform(X_train)
X_test_s = scaler.transform(X_test)
# Datasets
ds_train = SequenceDataset(
X=X_train_s,
window_size=window_size,
stride=stride,
labels=None,
label_strategy='last',
)
ds_test = SequenceDataset(
X=X_test_s,
window_size=window_size,
stride=stride,
labels=y_test,
label_strategy=label_strategy,
)
# Loaders
train_loader = DataLoader(
ds_train,
batch_size=batch_size,
shuffle=True,
num_workers=num_workers,
pin_memory=pin_memory,
drop_last=drop_last,
)
test_loader = DataLoader(
ds_test,
batch_size=batch_size,
shuffle=False,
num_workers=num_workers,
pin_memory=pin_memory,
drop_last=drop_last,
)
meta = dict(
feature_names=feature_names,
scaler=scaler,
min_=scaler.min_,
max_=scaler.max_,
train_rows=len(X_train),
test_rows=len(X_test),
window_size=window_size,
stride=stride,
label_strategy=label_strategy,
)
return train_loader, test_loader, meta
# -----------------------------
# CLI (optional quick check)
# -----------------------------
if __name__ == "__main__":
"""
Hardcoded one-click mode.
Edit the constants below to your own file paths and settings, then just run this file.
"""
import sys
# ====== EDIT THESE PATHS ======
TRAIN_CSV = r"E:\csv\features\ad1.csv" # 训练集 CSV(GBK),特征从第5列起
TEST_CSV = r"E:\csv\features\ad3-test4.csv" # 测试集 CSV(GBK),特征从第5列起,最后一列为标签
# ====== BASIC SETTINGS ======
WINDOW_SIZE = 64 # 滑窗长度
BATCH_SIZE = 32 # 批大小
STRIDE = 1 # 滑动步长
LABEL_STRATEGY = 'last' # 'last' | 'any' | 'majority'
ENCODING = 'gbk'
NUM_WORKERS = 0
PIN_MEMORY = False
DROP_LAST = False
# Build loaders (no CLI / no GUI)
try:
tr_loader, te_loader, meta = build_dataloaders(
train_csv=TRAIN_CSV,
test_csv=TEST_CSV,
window_size=WINDOW_SIZE,
batch_size=BATCH_SIZE,
stride=STRIDE,
label_strategy=LABEL_STRATEGY,
encoding=ENCODING,
num_workers=NUM_WORKERS,
pin_memory=PIN_MEMORY,
drop_last=DROP_LAST,
)
except Exception as e:
print("[ERROR] 加载数据失败:", e)
sys.exit(1)
print("\n✓ 数据加载完成(硬编码模式)")
print(f"训练 CSV: {TRAIN_CSV}")
print(f"测试 CSV: {TEST_CSV}")
print(f"窗口大小: {WINDOW_SIZE} | 批大小: {BATCH_SIZE} | 步长: {STRIDE}")
print(f"特征维度: {len(meta['feature_names'])} | 训练行数: {meta['train_rows']} | 测试行数: {meta['test_rows']}")
# Quick shape sanity checks
try:
xb = next(iter(tr_loader))
print("训练批次 shape:", tuple(xb.shape))
except StopIteration:
print("训练集窗口数为 0,请检查窗口大小/步长与数据行数是否匹配。")
print("Train windows:", len(tr_loader.dataset))
print("Test windows:", len(te_loader.dataset))
try:
batch = next(iter(te_loader))
if isinstance(batch, (list, tuple)) and len(batch) == 2:
xt, yt = batch
print("测试批次 x shape:", tuple(xt.shape), ", y shape:", tuple(yt.shape))
else:
print("测试批次 shape:", tuple(batch.shape))
except StopIteration:
print("测试集窗口数为 0,请检查窗口大小/步长与数据行数是否匹配。")