-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathddp_train_stella.py
More file actions
287 lines (242 loc) · 11.7 KB
/
Copy pathddp_train_stella.py
File metadata and controls
287 lines (242 loc) · 11.7 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
import torch
import torch.nn.functional as F
import torch.distributed as dist
from torch.nn.parallel import DistributedDataParallel as DDP
from torch.utils.data import DataLoader, DistributedSampler
from transformers import AutoTokenizer, AutoModel
from torch.optim import AdamW
from datasets import load_dataset
from tqdm import tqdm
import os
# LAT_MIN, LAT_MAX = -90.0, 90.0
# LON_MIN, LON_MAX = -180.0, 180.0
LAT_MIN = 20.0
LAT_MAX = 50.0
LON_MIN = -130.0
LON_MAX = -55.0
def filter_dataset(dataset, min_lat=LAT_MIN, max_lat=LAT_MAX, min_lng=LON_MIN, max_lng=LON_MAX):
"""
Filter the dataset to include only samples within the given latitude and longitude bounds.
Args:
dataset: The dataset to filter.
min_lat: Minimum latitude value.
max_lat: Maximum latitude value.
min_lng: Minimum longitude value.
max_lng: Maximum longitude value.
Returns:
Filtered dataset.
"""
return dataset.filter(lambda item:
min_lat <= item['latitude'] <= max_lat and
min_lng <= item['longitude'] <= max_lng
)
# Define the model
class TextToGeoModel(torch.nn.Module):
def __init__(self, embedding_model_path, embedding_dim=1024):
super(TextToGeoModel, self).__init__()
self.embedding_dim = embedding_dim
self.tokenizer = AutoTokenizer.from_pretrained(embedding_model_path, trust_remote_code=True)
self.embedding_model = AutoModel.from_pretrained(embedding_model_path, trust_remote_code=True)
self.fc = torch.nn.Linear(self.embedding_model.config.hidden_size, 2)
# self.vector_linear = torch.nn.Linear(
# in_features=self.embedding_model.config.hidden_size, out_features=embedding_dim
# )
# self.vector_linear_dict = {
# k.replace("linear.", ""): v
# for k, v in torch.load(os.path.join(embedding_model_path, "2_Dense_{}/pytorch_model.bin".format(embedding_dim))).items()
# }
# self.vector_linear.load_state_dict(self.vector_linear_dict)
# self.vector_linear.cuda()
# self.fc = torch.nn.Linear(embedding_dim, 2)
# Freeze the embedding model parameters if necessary
for param in self.embedding_model.parameters():
param.requires_grad = False
def mean_pooling(self, hidden_state, attention_mask):
hidden_state = hidden_state.masked_fill(~attention_mask[..., None].bool(), 0.0)
return hidden_state.sum(dim=1) / attention_mask.sum(dim=1)[..., None]
def forward(self, texts):
if isinstance(texts, str):
texts = [texts]
encoded_input = self.tokenizer(texts, padding=True, truncation=True, max_length=512, return_tensors='pt')
encoded_input = {key: val.to(next(self.embedding_model.parameters()).device) for key, val in encoded_input.items()}
attention_mask = encoded_input['attention_mask']
last_hidden_state = self.embedding_model(**encoded_input).last_hidden_state
pooled_output = self.mean_pooling(last_hidden_state, attention_mask)
# reduced_output = self.vector_linear(pooled_output) # Apply the additional linear layer
embeddings = F.normalize(pooled_output, p=2, dim=-1)
return self.fc(embeddings)
def freeze_embedding_model(self):
for param in self.embedding_model.parameters():
param.requires_grad = False
# Dataset collate function
# def collate_fn(batch):
# texts = [item['text'][:100] for item in batch]
# targets = torch.tensor([[item['latitude'], item['longitude']] for item in batch], dtype=torch.float32)
# return texts, targets
# Dataset collate function with normalization
def collate_fn(batch):
texts = [item['text'] for item in batch]
targets = torch.tensor([
[
(item['latitude'] - LAT_MIN) / (LAT_MAX - LAT_MIN), # Normalize latitude
(item['longitude'] - LON_MIN) / (LON_MAX - LON_MIN) # Normalize longitude
]
for item in batch
], dtype=torch.float32)
return texts, targets
# Function to save model checkpoint
def save_checkpoint(model, optimizer, epoch, loss, checkpoint_dir="checkpoints", filename="best_checkpoint.pth"):
os.makedirs(checkpoint_dir, exist_ok=True)
checkpoint_path = os.path.join(checkpoint_dir, filename)
torch.save({
'epoch': epoch,
'model_state_dict': model.state_dict(),
'optimizer_state_dict': optimizer.state_dict(),
'loss': loss,
}, checkpoint_path)
print(f"Checkpoint saved: {checkpoint_path}")
# def haversine_distance_loss(preds, targets):
# """
# Compute the mean Haversine distance between predicted and target coordinates.
# Args:
# preds: Predicted coordinates of shape (N, 2) [latitude, longitude] in degrees.
# targets: Ground truth coordinates of shape (N, 2) [latitude, longitude] in degrees.
# Returns:
# Mean Haversine distance in kilometers.
# """
# R = 6371.0 # Earth's radius in kilometers
# # Convert latitude and longitude from degrees to radians
# lat_pred, lon_pred = torch.deg2rad(preds[:, 0]), torch.deg2rad(preds[:, 1])
# lat_target, lon_target = torch.deg2rad(targets[:, 0]), torch.deg2rad(targets[:, 1])
# # Compute differences
# delta_lat = lat_target - lat_pred
# delta_lon = lon_target - lon_pred
# # Haversine formula
# a = torch.sin(delta_lat / 2)**2 + torch.cos(lat_pred) * torch.cos(lat_target) * torch.sin(delta_lon / 2)**2
# c = 2 * torch.atan2(torch.sqrt(a), torch.sqrt(1 - a))
# haversine_distance = R * c
# return haversine_distance.mean()
def haversine_distance_loss(preds, targets):
"""
Compute the mean Haversine distance between predicted and target coordinates.
Args:
preds: Normalized predicted coordinates (N, 2).
targets: Normalized ground truth coordinates (N, 2).
Returns:
Mean Haversine distance in kilometers.
"""
# Denormalize predictions and targets
preds[:, 0] = preds[:, 0] * (LAT_MAX - LAT_MIN) + LAT_MIN # Denormalize latitude
preds[:, 1] = preds[:, 1] * (LON_MAX - LON_MIN) + LON_MIN # Denormalize longitude
targets[:, 0] = targets[:, 0] * (LAT_MAX - LAT_MIN) + LAT_MIN
targets[:, 1] = targets[:, 1] * (LON_MAX - LON_MIN) + LON_MIN
# Haversine distance computation (same as before)
R = 6371.0 # Earth's radius in kilometers
lat_pred, lon_pred = torch.deg2rad(preds[:, 0]), torch.deg2rad(preds[:, 1])
lat_target, lon_target = torch.deg2rad(targets[:, 0]), torch.deg2rad(targets[:, 1])
delta_lat = lat_target - lat_pred
delta_lon = lon_target - lon_pred
a = torch.sin(delta_lat / 2)**2 + torch.cos(lat_pred) * torch.cos(lat_target) * torch.sin(delta_lon / 2)**2
c = 2 * torch.atan2(torch.sqrt(a), torch.sqrt(1 - a))
haversine_distance = R * c
return haversine_distance.mean()
# Denormalize outputs before returning
def denormalize_coordinates(preds):
preds[:, 0] = preds[:, 0] * (LAT_MAX - LAT_MIN) + LAT_MIN
preds[:, 1] = preds[:, 1] * (LON_MAX - LON_MIN) + LON_MIN
return preds
# Training script with checkpoint saving and early exit
def main(local_rank, num_epochs=5, batch_size=512, patience=3):
# Initialize the process group
dist.init_process_group(backend="nccl", init_method="env://")
torch.cuda.set_device(local_rank)
device = torch.device("cuda", local_rank)
# Load dataset
dataset_name = 'ai-practicum-group/post2geo-dataset'
post2geo_dataset = load_dataset(dataset_name, token=os.environ.get('HF_TOKEN'))
# Filter train, validation, and test splits
post2geo_dataset['train'] = filter_dataset(post2geo_dataset['train'])
post2geo_dataset['validation'] = filter_dataset(post2geo_dataset['validation'])
post2geo_dataset['test'] = filter_dataset(post2geo_dataset['test'])
print("Data After Filtering:\n", post2geo_dataset)
# Dataloader with DistributedSampler
train_sampler = DistributedSampler(post2geo_dataset['train'], shuffle=True)
train_dataloader = DataLoader(
post2geo_dataset['train'], batch_size=batch_size, sampler=train_sampler, collate_fn=collate_fn, pin_memory=True, num_workers=64
)
val_sampler = DistributedSampler(post2geo_dataset['validation'], shuffle=False)
val_dataloader = DataLoader(
post2geo_dataset['validation'], batch_size=batch_size, sampler=val_sampler, collate_fn=collate_fn, pin_memory=True, num_workers=64
)
test_sampler = DistributedSampler(post2geo_dataset['test'], shuffle=False)
test_dataloader = DataLoader(
post2geo_dataset['test'], batch_size=batch_size, sampler=test_sampler, collate_fn=collate_fn, pin_memory=True, num_workers=64
)
# Model and optimizer
model = TextToGeoModel(embedding_model_path="dunzhang/stella_en_400M_v5", embedding_dim=1024).to(device)
model = DDP(model, device_ids=[local_rank], find_unused_parameters=False)
optimizer = AdamW(model.parameters(), lr=5e-5)
criterion = haversine_distance_loss
# Variables for tracking the best loss and early stopping
best_loss = float('inf')
best_epoch = -1
no_improvement_count = 0 # Tracks epochs without improvement
# Training loop
for epoch in range(num_epochs):
model.train()
train_sampler.set_epoch(epoch)
train_loss = 0.0
for texts, targets in tqdm(train_dataloader, desc=f"Epoch {epoch+1}/{num_epochs}", disable=(local_rank != 0)):
texts = [text for text in texts]
targets = targets.to(device)
optimizer.zero_grad()
outputs = model(texts)
loss = criterion(outputs, targets)
loss.backward()
optimizer.step()
train_loss += loss.item()
if local_rank == 0:
print(f"Epoch {epoch+1} - Training Loss: {train_loss/len(train_dataloader)} km")
# Evaluation
if local_rank == 0:
model.eval()
val_loss = 0.0
with torch.no_grad():
for texts, targets in tqdm(val_dataloader, desc="Validation"):
texts = [text for text in texts]
targets = targets.to(device)
outputs = model(texts)
outputs = denormalize_coordinates(outputs)
loss = criterion(outputs, targets)
val_loss += loss.item()
avg_val_loss = val_loss / len(val_dataloader)
print(f"Epoch {epoch+1} - Validation Loss: {avg_val_loss} km")
test_loss = 0.0
with torch.no_grad():
for texts, targets in tqdm(test_dataloader, desc="Validation"):
texts = [text for text in texts]
targets = targets.to(device)
outputs = model(texts)
outputs = denormalize_coordinates(outputs)
loss = criterion(outputs, targets)
test_loss += loss.item()
avg_test_loss = test_loss / len(test_dataloader)
print(f"Epoch {epoch+1} - Test Loss: {avg_test_loss} km")
# Save checkpoint if this is the best model so far
if avg_val_loss < best_loss:
best_loss = avg_val_loss
best_epoch = epoch
no_improvement_count = 0 # Reset no improvement count
save_checkpoint(model, optimizer, epoch, best_loss)
else:
no_improvement_count += 1 # Increment no improvement count
# Check for early exit
if no_improvement_count >= patience:
print(f"Early stopping triggered after {patience} epochs without improvement.")
break
# Final cleanup
dist.destroy_process_group()
if __name__ == "__main__":
# Entry point for torchrun
local_rank = int(os.environ["LOCAL_RANK"])
main(local_rank)