Merge pull request #27 from real-lhj/main

add AnomalyTransformer
This commit is contained in:
Hongzuo Xu 2023-09-03 20:51:31 +08:00 committed by GitHub
commit 71cac7ed3b
No known key found for this signature in database
GPG Key ID: 4AEE18F83AFDEB23
6 changed files with 579 additions and 14 deletions

View File

@ -5,6 +5,7 @@ from .tranad import TranAD
from .usad import USAD
from .couta import COUTA
from .tcned import TcnED
from .anomalytransformer import AnomalyTransformer
# weakly-supervised
from .dsad import DeepSADTS
@ -13,4 +14,4 @@ from .prenet import PReNetTS
__all__ = ['DeepIsolationForestTS', 'DeepSVDDTS', 'TranAD', 'USAD', 'COUTA',
'DeepSADTS', 'DevNetTS', 'PReNetTS']
'DeepSADTS', 'DevNetTS', 'PReNetTS', 'AnomalyTransformer']

View File

@ -0,0 +1,407 @@
import torch
import torch.nn as nn
import torch.nn.functional as F
import numpy as np
from torch.utils.data import DataLoader
import math
from deepod.utils.utility import get_sub_seqs
from deepod.core.base_model import BaseDeepAD
def my_kl_loss(p, q):
res = p * (torch.log(p + 0.0001) - torch.log(q + 0.0001))
return torch.mean(torch.sum(res, dim=-1), dim=1)
class AnomalyTransformer(BaseDeepAD):
def __init__(self, seq_len=100, stride=1, lr=0.0001, epochs=10, batch_size=32,
epoch_steps=20, prt_steps=1, device='cuda',
k=3, input_c=25, output_c=25, anomaly_ratio=1,
verbose=2, random_state=42):
super(AnomalyTransformer, self).__init__(
model_name='AnomalyTransformer', data_type='ts', epochs=epochs, batch_size=batch_size, lr=lr,
seq_len=seq_len, stride=stride,
epoch_steps=epoch_steps, prt_steps=prt_steps, device=device,
verbose=verbose, random_state=random_state
)
self.k = k
self.input_c = input_c
self.output_c = output_c
self.anomaly_ratio = anomaly_ratio
def fit(self, X, y=None):
self.n_features = X.shape[1]
train_seqs = get_sub_seqs(X, seq_len=self.seq_len, stride=self.stride)
self.model = AnomalyTransformerModel(
win_size=self.seq_len,
enc_in=self.input_c,
c_out=self.output_c,
e_layers=3).to(self.device)
dataloader = DataLoader(train_seqs, batch_size=self.batch_size,
shuffle=True, pin_memory=True)
self.optimizer = torch.optim.AdamW(self.model.parameters(), lr=self.lr, weight_decay=1e-5)
self.scheduler = torch.optim.lr_scheduler.StepLR(self.optimizer, step_size=5, gamma=0.5)
self.model.train()
for e in range(self.epochs):
loss = self.training(dataloader)
print(f'Epoch {e + 1},\t L1 = {loss}')
self.decision_scores_ = self.decision_function(X)
self.labels_ = self._process_decision_scores() # in base model
return
def decision_function(self, X, return_rep=False):
seqs = get_sub_seqs(X, seq_len=self.seq_len, stride=1)
dataloader = DataLoader(seqs, batch_size=self.batch_size,
shuffle=False, drop_last=False)
self.model.eval()
loss, _ = self.inference(dataloader) # (n,d)
loss_final = np.mean(loss, axis=1) # (n,)
padding_list = np.zeros([X.shape[0] - loss.shape[0], loss.shape[1]])
loss_pad = np.concatenate([padding_list, loss], axis=0)
loss_final_pad = np.hstack([0 * np.ones(X.shape[0] - loss_final.shape[0]), loss_final])
return loss_final_pad
def training(self, dataloader):
criterion = nn.MSELoss()
loss_list = []
for ii, batch_x in enumerate(dataloader):
self.optimizer.zero_grad()
input = batch_x.float().to(self.device)
output, series, prior, _ = self.model(input)
# calculate Association discrepancy
series_loss = 0.0
prior_loss = 0.0
for u in range(len(prior)):
series_loss += (torch.mean(my_kl_loss(series[u], (
prior[u] / torch.unsqueeze(torch.sum(prior[u], dim=-1), dim=-1).repeat(1, 1, 1,
self.seq_len)).detach())) + torch.mean(
my_kl_loss((prior[u] / torch.unsqueeze(torch.sum(prior[u], dim=-1), dim=-1).repeat(1, 1, 1,
self.seq_len)).detach(),
series[u])))
prior_loss += (torch.mean(my_kl_loss(
(prior[u] / torch.unsqueeze(torch.sum(prior[u], dim=-1), dim=-1).repeat(1, 1, 1,
self.seq_len)),
series[u].detach())) + torch.mean(
my_kl_loss(series[u].detach(), (
prior[u] / torch.unsqueeze(torch.sum(prior[u], dim=-1), dim=-1).repeat(1, 1, 1,
self.seq_len)))))
series_loss = series_loss / len(prior)
prior_loss = prior_loss / len(prior)
rec_loss = criterion(output, input) # compute loss
loss_list.append((rec_loss - self.k * series_loss).item())
loss1 = rec_loss - self.k * series_loss
loss2 = rec_loss + self.k * prior_loss
# Minimax strategy
loss1.backward(retain_graph=True)
loss2.backward()
self.optimizer.step()
if self.epoch_steps != -1:
if ii > self.epoch_steps:
break
self.scheduler.step()
return np.average(loss_list)
def inference(self, dataloader):
criterion = nn.MSELoss(reduce=False)
temperature = 50
attens_energy = []
preds = []
for input_data in dataloader: # test_set
input = input_data.float().to(self.device)
output, series, prior, _ = self.model(input)
loss = torch.mean(criterion(input, output), dim=-1)
series_loss = 0.0
prior_loss = 0.0
for u in range(len(prior)):
if u == 0:
series_loss = my_kl_loss(series[u], (
prior[u] / torch.unsqueeze(torch.sum(prior[u], dim=-1), dim=-1).repeat(1, 1, 1,
self.seq_len)).detach()) * temperature
prior_loss = my_kl_loss(
(prior[u] / torch.unsqueeze(torch.sum(prior[u], dim=-1), dim=-1).repeat(1, 1, 1,
self.seq_len)),
series[u].detach()) * temperature
else:
series_loss += my_kl_loss(series[u], (
prior[u] / torch.unsqueeze(torch.sum(prior[u], dim=-1), dim=-1).repeat(1, 1, 1,
self.seq_len)).detach()) * temperature
prior_loss += my_kl_loss(
(prior[u] / torch.unsqueeze(torch.sum(prior[u], dim=-1), dim=-1).repeat(1, 1, 1,
self.seq_len)),
series[u].detach()) * temperature
metric = torch.softmax((-series_loss - prior_loss), dim=-1)
cri = metric * loss
cri = cri.detach().cpu().numpy()
attens_energy.append(cri)
attens_energy = np.concatenate(attens_energy, axis=0) # anomaly scores
test_energy = np.array(attens_energy) # anomaly scores
return test_energy, preds # (n,d)
def training_forward(self, batch_x, net, criterion):
"""define forward step in training"""
return
def inference_forward(self, batch_x, net, criterion):
"""define forward step in inference"""
return
def training_prepare(self, X, y):
"""define train_loader, net, and criterion"""
return
def inference_prepare(self, X):
"""define test_loader"""
return
# Proposed Model
class AnomalyTransformerModel(nn.Module):
def __init__(self, win_size, enc_in, c_out, d_model=512, n_heads=8, e_layers=3, d_ff=512,
dropout=0.0, activation='gelu', output_attention=True):
super(AnomalyTransformerModel, self).__init__()
self.output_attention = output_attention
# Encoding
self.embedding = DataEmbedding(enc_in, d_model, dropout)
# Encoder
self.encoder = Encoder(
[
EncoderLayer(
AttentionLayer(
AnomalyAttention(win_size, False, attention_dropout=dropout, output_attention=output_attention),
d_model, n_heads),
d_model,
d_ff,
dropout=dropout,
activation=activation
) for l in range(e_layers)
],
norm_layer=torch.nn.LayerNorm(d_model)
)
self.projection = nn.Linear(d_model, c_out, bias=True)
def forward(self, x):
enc_out = self.embedding(x)
enc_out, series, prior, sigmas = self.encoder(enc_out)
enc_out = self.projection(enc_out)
if self.output_attention:
return enc_out, series, prior, sigmas
else:
return enc_out # [B, L, D]
class EncoderLayer(nn.Module):
def __init__(self, attention, d_model, d_ff=None, dropout=0.1, activation="relu"):
super(EncoderLayer, self).__init__()
d_ff = d_ff or 4 * d_model
self.attention = attention
self.conv1 = nn.Conv1d(in_channels=d_model, out_channels=d_ff, kernel_size=1)
self.conv2 = nn.Conv1d(in_channels=d_ff, out_channels=d_model, kernel_size=1)
self.norm1 = nn.LayerNorm(d_model)
self.norm2 = nn.LayerNorm(d_model)
self.dropout = nn.Dropout(dropout)
self.activation = F.relu if activation == "relu" else F.gelu
def forward(self, x, attn_mask=None):
new_x, attn, mask, sigma = self.attention(
x, x, x,
attn_mask=attn_mask
)
x = x + self.dropout(new_x)
y = x = self.norm1(x)
y = self.dropout(self.activation(self.conv1(y.transpose(-1, 1))))
y = self.dropout(self.conv2(y).transpose(-1, 1))
return self.norm2(x + y), attn, mask, sigma
class Encoder(nn.Module):
def __init__(self, attn_layers, norm_layer=None):
super(Encoder, self).__init__()
self.attn_layers = nn.ModuleList(attn_layers)
self.norm = norm_layer
def forward(self, x, attn_mask=None):
# x [B, L, D]
series_list = []
prior_list = []
sigma_list = []
for attn_layer in self.attn_layers:
x, series, prior, sigma = attn_layer(x, attn_mask=attn_mask)
series_list.append(series)
prior_list.append(prior)
sigma_list.append(sigma)
if self.norm is not None:
x = self.norm(x)
return x, series_list, prior_list, sigma_list
class DataEmbedding(nn.Module):
def __init__(self, c_in, d_model, dropout=0.0):
super(DataEmbedding, self).__init__()
self.value_embedding = TokenEmbedding(c_in=c_in, d_model=d_model)
self.position_embedding = PositionalEmbedding(d_model=d_model)
self.dropout = nn.Dropout(p=dropout)
def forward(self, x):
x = self.value_embedding(x) + self.position_embedding(x)
return self.dropout(x)
class TokenEmbedding(nn.Module):
def __init__(self, c_in, d_model):
super(TokenEmbedding, self).__init__()
padding = 1 if torch.__version__ >= '1.5.0' else 2
self.tokenConv = nn.Conv1d(in_channels=c_in, out_channels=d_model,
kernel_size=3, padding=padding, padding_mode='circular', bias=False)
for m in self.modules():
if isinstance(m, nn.Conv1d):
nn.init.kaiming_normal_(m.weight, mode='fan_in', nonlinearity='leaky_relu')
def forward(self, x):
x = self.tokenConv(x.permute(0, 2, 1)).transpose(1, 2)
return x
class PositionalEmbedding(nn.Module):
def __init__(self, d_model, max_len=5000):
super(PositionalEmbedding, self).__init__()
# Compute the positional encodings once in log space.
pe = torch.zeros(max_len, d_model).float()
pe.require_grad = False
position = torch.arange(0, max_len).float().unsqueeze(1)
div_term = (torch.arange(0, d_model, 2).float() * -(math.log(10000.0) / d_model)).exp()
pe[:, 0::2] = torch.sin(position * div_term)
pe[:, 1::2] = torch.cos(position * div_term)
pe = pe.unsqueeze(0)
self.register_buffer('pe', pe)
def forward(self, x):
return self.pe[:, :x.size(1)]
class TriangularCausalMask():
def __init__(self, B, L, device="cpu"):
mask_shape = [B, 1, L, L]
with torch.no_grad():
self._mask = torch.triu(torch.ones(mask_shape, dtype=torch.bool), diagonal=1).to(device)
@property
def mask(self):
return self._mask
class AnomalyAttention(nn.Module):
def __init__(self, win_size, mask_flag=True, scale=None, attention_dropout=0.0, output_attention=False):
super(AnomalyAttention, self).__init__()
self.scale = scale
self.mask_flag = mask_flag
self.output_attention = output_attention
self.dropout = nn.Dropout(attention_dropout)
window_size = win_size
self.distances = torch.zeros((window_size, window_size)).cuda()
for i in range(window_size):
for j in range(window_size):
self.distances[i][j] = abs(i - j)
def forward(self, queries, keys, values, sigma, attn_mask):
B, L, H, E = queries.shape
_, S, _, D = values.shape
scale = self.scale or 1. / math.sqrt(E)
scores = torch.einsum("blhe,bshe->bhls", queries, keys)
if self.mask_flag:
if attn_mask is None:
attn_mask = TriangularCausalMask(B, L, device=queries.device)
scores.masked_fill_(attn_mask.mask, -np.inf)
attn = scale * scores
sigma = sigma.transpose(1, 2) # B L H -> B H L
window_size = attn.shape[-1]
sigma = torch.sigmoid(sigma * 5) + 1e-5
sigma = torch.pow(3, sigma) - 1
sigma = sigma.unsqueeze(-1).repeat(1, 1, 1, window_size) # B H L L
prior = self.distances.unsqueeze(0).unsqueeze(0).repeat(sigma.shape[0], sigma.shape[1], 1, 1).cuda()
prior = 1.0 / (math.sqrt(2 * math.pi) * sigma) * torch.exp(-prior ** 2 / 2 / (sigma ** 2))
series = self.dropout(torch.softmax(attn, dim=-1))
V = torch.einsum("bhls,bshd->blhd", series, values)
if self.output_attention:
return (V.contiguous(), series, prior, sigma)
else:
return (V.contiguous(), None)
class AttentionLayer(nn.Module):
def __init__(self, attention, d_model, n_heads, d_keys=None,
d_values=None):
super(AttentionLayer, self).__init__()
d_keys = d_keys or (d_model // n_heads)
d_values = d_values or (d_model // n_heads)
self.norm = nn.LayerNorm(d_model)
self.inner_attention = attention
self.query_projection = nn.Linear(d_model,
d_keys * n_heads)
self.key_projection = nn.Linear(d_model,
d_keys * n_heads)
self.value_projection = nn.Linear(d_model,
d_values * n_heads)
self.sigma_projection = nn.Linear(d_model,
n_heads)
self.out_projection = nn.Linear(d_values * n_heads, d_model)
self.n_heads = n_heads
def forward(self, queries, keys, values, attn_mask):
B, L, _ = queries.shape
_, S, _ = keys.shape
H = self.n_heads
x = queries
queries = self.query_projection(queries).view(B, L, H, -1)
keys = self.key_projection(keys).view(B, S, H, -1)
values = self.value_projection(values).view(B, S, H, -1)
sigma = self.sigma_projection(x).view(B, L, H)
out, series, prior, sigma = self.inner_attention(
queries,
keys,
values,
sigma,
attn_mask
)
out = out.view(B, L, -1)
return self.out_projection(out), series, prior, sigma

View File

@ -57,12 +57,12 @@ class TranAD(BaseDeepAD):
shuffle=False, drop_last=False)
self.model.eval()
loss, _ = self.inference(dataloader)
loss_final = np.mean(loss, axis=1)
loss, _ = self.inference(dataloader) # (8611,d)
loss_final = np.mean(loss, axis=1) # (n,)
padding_list = np.zeros([X.shape[0]-loss.shape[0], loss.shape[1]])
loss_pad = np.concatenate([padding_list, loss], axis=0)
loss_final_pad = np.hstack([0 * np.ones(X.shape[0] - loss_final.shape[0]), loss_final])
loss_final_pad = np.hstack([0 * np.ones(X.shape[0] - loss_final.shape[0]), loss_final]) # (8640,)
return loss_final_pad
@ -73,15 +73,15 @@ class TranAD(BaseDeepAD):
l1s, l2s = [], []
for ii, batch_x in enumerate(dataloader):
local_bs = batch_x.shape[0]
window = batch_x.permute(1, 0, 2)
elem = window[-1, :, :].view(1, local_bs, self.n_features)
local_bs = batch_x.shape[0] #(1283019)
window = batch_x.permute(1, 0, 2) # (30, 128, 19)
elem = window[-1, :, :].view(1, local_bs, self.n_features) #(1, 128, 19)
window = window.float().to(self.device)
elem = elem.float().to(self.device)
z = self.model(window, elem)
l1 = (1/n) * criterion(z[0], elem) + (1-1/n) * criterion(z[1], elem)
l1 = (1/n) * criterion(z[0], elem) + (1-1/n) * criterion(z[1], elem) #(1, 128, 19)
l1s.append(torch.mean(l1).item())
loss = torch.mean(l1)

View File

@ -0,0 +1,141 @@
# -*- coding: utf-8 -*-
from __future__ import division
from __future__ import print_function
import os
import sys
import unittest
# noinspection PyProtectedMember
from numpy.testing import assert_equal
from sklearn.metrics import roc_auc_score
import torch
import pandas as pd
# temporary solution for relative imports in case pyod is not installed
# if deepod is installed, no need to use the following line
sys.path.append(os.path.abspath(os.path.join(os.path.dirname(__file__), '..')))
from deepod.models.time_series.anomalytransformer import AnomalyTransformer
class TestAnomalyTransformer(unittest.TestCase):
def setUp(self):
train_file = 'data/omi-1/omi-1_train.csv'
test_file = 'data/omi-1/omi-1_test.csv'
train_df = pd.read_csv(train_file, sep=',', index_col=0)
test_df = pd.read_csv(test_file, index_col=0)
y = test_df['label'].values
train_df, test_df = train_df.drop('label', axis=1), test_df.drop('label', axis=1)
self.Xts_train = train_df.values
self.Xts_test = test_df.values
self.yts_test = y
device = 'cuda' if torch.cuda.is_available() else 'cpu'
self.clf = AnomalyTransformer(seq_len=100, stride=1, epochs=10,
batch_size=32, k=3, input_c=19, output_c=19, anomaly_ratio=1, lr=1e-4,
device=device, random_state=42)
self.clf.fit(self.Xts_train)
def test_parameters(self):
assert (hasattr(self.clf, 'decision_scores_') and
self.clf.decision_scores_ is not None)
assert (hasattr(self.clf, 'labels_') and
self.clf.labels_ is not None)
assert (hasattr(self.clf, 'threshold_') and
self.clf.threshold_ is not None)
def test_train_scores(self):
assert_equal(len(self.clf.decision_scores_), self.Xts_train.shape[0])
def test_prediction_scores(self):
pred_scores = self.clf.decision_function(self.Xts_test)
assert_equal(pred_scores.shape[0], self.Xts_test.shape[0])
def test_prediction_labels(self):
pred_labels = self.clf.predict(self.Xts_test)
assert_equal(pred_labels.shape, self.yts_test.shape)
# def test_prediction_proba(self):
# pred_proba = self.clf.predict_proba(self.X_test)
# assert (pred_proba.min() >= 0)
# assert (pred_proba.max() <= 1)
#
# def test_prediction_proba_linear(self):
# pred_proba = self.clf.predict_proba(self.X_test, method='linear')
# assert (pred_proba.min() >= 0)
# assert (pred_proba.max() <= 1)
#
# def test_prediction_proba_unify(self):
# pred_proba = self.clf.predict_proba(self.X_test, method='unify')
# assert (pred_proba.min() >= 0)
# assert (pred_proba.max() <= 1)
#
# def test_prediction_proba_parameter(self):
# with assert_raises(ValueError):
# self.clf.predict_proba(self.X_test, method='something')
def test_prediction_labels_confidence(self):
pred_labels, confidence = self.clf.predict(self.Xts_test, return_confidence=True)
assert_equal(pred_labels.shape, self.yts_test.shape)
assert_equal(confidence.shape, self.yts_test.shape)
assert (confidence.min() >= 0)
assert (confidence.max() <= 1)
# def test_prediction_proba_linear_confidence(self):
# pred_proba, confidence = self.clf.predict_proba(self.X_test,
# method='linear',
# return_confidence=True)
# assert (pred_proba.min() >= 0)
# assert (pred_proba.max() <= 1)
#
# assert_equal(confidence.shape, self.y_test.shape)
# assert (confidence.min() >= 0)
# assert (confidence.max() <= 1)
#
# def test_fit_predict(self):
# pred_labels = self.clf.fit_predict(self.X_train)
# assert_equal(pred_labels.shape, self.y_train.shape)
#
# def test_fit_predict_score(self):
# self.clf.fit_predict_score(self.X_test, self.y_test)
# self.clf.fit_predict_score(self.X_test, self.y_test,
# scoring='roc_auc_score')
# self.clf.fit_predict_score(self.X_test, self.y_test,
# scoring='prc_n_score')
# with assert_raises(NotImplementedError):
# self.clf.fit_predict_score(self.X_test, self.y_test,
# scoring='something')
#
# def test_predict_rank(self):
# pred_socres = self.clf.decision_function(self.X_test)
# pred_ranks = self.clf._predict_rank(self.X_test)
#
# # assert the order is reserved
# assert_allclose(rankdata(pred_ranks), rankdata(pred_socres), atol=3)
# assert_array_less(pred_ranks, self.X_train.shape[0] + 1)
# assert_array_less(-0.1, pred_ranks)
#
# def test_predict_rank_normalized(self):
# pred_socres = self.clf.decision_function(self.X_test)
# pred_ranks = self.clf._predict_rank(self.X_test, normalized=True)
#
# # assert the order is reserved
# assert_allclose(rankdata(pred_ranks), rankdata(pred_socres), atol=3)
# assert_array_less(pred_ranks, 1.01)
# assert_array_less(-0.1, pred_ranks)
# def test_plot(self):
# os, cutoff1, cutoff2 = self.clf.explain_outlier(ind=1)
# assert_array_less(0, os)
# def test_model_clone(self):
# clone_clf = clone(self.clf)
def tearDown(self):
pass
if __name__ == '__main__':
unittest.main()

View File

@ -13,11 +13,14 @@ def get_model_class(name):
return USAD
elif name == 'couta':
return COUTA
elif name == 'anomalytransformer':
return AnomalyTransformer
# elif name == 'seqregad':
# return SeqRegAD
def get_additional_configs(name):
def get_additional_configs(model_name, ds_name):
n_dims = {"SMAP": 25, "SMD": 38, "ASD": 19}
config_dict = {
# 'seqregad': {
# 'seq_len_lst': [10, 30, 50],
@ -67,10 +70,20 @@ def get_additional_configs(name):
'lr': 1e-4,
'epochs': 20,
'batch_size': 64,
},
'anomalytransformer': {
'lr': 1e-4,
'epochs': 10,
'batch_size': 32,
'k': 3,
'input_c': n_dims[ds_name],
'output_c': n_dims[ds_name],
'anomaly_ratio': 1
}
}
try:
return config_dict[name]
return config_dict[model_name]
except KeyError:
return {}

View File

@ -15,20 +15,23 @@ import configs
dataset_root = f'/home/{getpass.getuser()}/dataset/5-TSdata/_processed_data/'
parser = argparse.ArgumentParser()
parser.add_argument("--runs", type=int, default=5,
help="how many times we repeat the experiments to obtain the average performance")
parser.add_argument("--output_dir", type=str, default='@records/',
help="the output file path")
parser.add_argument("--dataset", type=str,
default='ASD,SMAP,MSL',
)
parser.add_argument("--entities", type=str,
default='FULL',
help='FULL represents all the csv file in the folder, or a list of entity names split by comma'
)
parser.add_argument("--entity_combined", type=int, default=1)
parser.add_argument("--model", type=str, default='tranad', help="")
parser.add_argument("--model", type=str, default='anomalytransformer', help="") # anomalytransformer
parser.add_argument('--silent_header', action='store_true')
parser.add_argument("--flag", type=str, default='')
@ -40,7 +43,7 @@ parser.add_argument('--stride', type=int, default=10)
args = parser.parse_args()
model_class = configs.get_model_class(args.model)
model_configs = configs.get_additional_configs(args.model)
model_configs = configs.get_additional_configs(args.model, args.dataset)
model_configs['seq_len'] = args.seq_len
model_configs['stride'] = args.stride
@ -88,8 +91,8 @@ for dataset in dataset_name_lst:
t1 = time.time()
clf = model_class(**model_configs, random_state=42+i)
clf.fit(train_data)
scores = clf.decision_function(test_data)
clf.fit(train_data) # 主要改这里
scores = clf.decision_function(test_data) # scores(n,) test_data(n,d)
t = time.time() - t1
eval_metrics = utils.get_metrics(labels, scores)