Compare commits

...
Author SHA1 Message Date
LysandreJik df2458c2c9 Pytorch first experiments 2019-10-09 11:02:14 -04:00
LysandreJik 111bf7ca07 Credit where it's due 2019-10-04 14:33:31 -04:00
LysandreJik 98cbe5c6fa Refactoring 2019-10-04 14:21:59 -04:00
LysandreJik 2445e642ca Working dataset 2019-10-03 11:26:35 -04:00
LysandreJik 0af5ef2090 WIP 2019-10-02 17:41:52 -04:00
LysandreJik d7d1e0523b Not using sharding 2019-09-30 18:42:02 -04:00
LysandreJik 234d91e3c3 Using fit 2019-09-29 13:12:38 -04:00
LysandreJik 28ecf31833 Initial TPU experiments 2019-09-29 12:53:18 -04:00
9 changed files with 545 additions and 42 deletions
+24
View File
@@ -0,0 +1,24 @@
import os
os.environ["TPU_IP_ADDRESS"] = "192.168.0.2"
os.environ["TPU_NAME"] = "node-1"
os.environ["XRT_TPU_CONFIG"] = "tpu_worker;0;192.168.0.2:8470"
import torch
import torch_xla
import torch_xla.core.xla_model as xm
device = xm.xla_device()
from transformers import GPT2LMHeadModel, GPT2Tokenizer
model = GPT2LMHeadModel.from_pretrained("gpt2")
tokenizer = GPT2Tokenizer.from_pretrained("gpt2")
sequence = "This runs on TPU"
input_ids = torch.tensor([tokenizer.encode(sequence)], device=device)
model.train().to(device)
print(input_ids)
+267
View File
@@ -0,0 +1,267 @@
# coding=utf-8
# Copyright 2018 The Open AI Team Authors and The HuggingFace Inc. team.
#
# Licensed under the Apache License, Version 2.0 (the "License");
# you may not use this file except in compliance with the License.
# You may obtain a copy of the License at
#
# http://www.apache.org/licenses/LICENSE-2.0
#
# Unless required by applicable law or agreed to in writing, software
# distributed under the License is distributed on an "AS IS" BASIS,
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
# See the License for the specific language governing permissions and
# limitations under the License.
""" === Under active development === Script to fine-tune GLUE on a TPU
Adapted from https://github.com/tensorflow/models
Especially https://github.com/tensorflow/models/blob/master/official/modeling/model_training_utils.py
"""
from transformers import TFBertForSequenceClassification, BertTokenizer, BertConfig
from tpu_utils import get_tpu
from tpu_dataset import create_dataset
import functools
import tensorflow as tf
import logging
import json
import math
def _get_input_iterator(input_fn, strategy):
"""Returns distributed dataset iterator."""
# When training with TPU pods, datasets needs to be cloned across
# workers. Since Dataset instance cannot be cloned in eager mode, we instead
# pass callable that returns a dataset.
input_data = input_fn()
if callable(input_data):
iterator = iter(
strategy.experimental_distribute_datasets_from_function(input_data))
else:
iterator = iter(strategy.experimental_distribute_dataset(input_data))
return iterator
def _steps_to_run(current_step, steps_per_epoch, steps_per_loop):
"""Calculates steps to run on device."""
if steps_per_loop <= 0:
raise ValueError('steps_per_loop should be positive integer.')
if steps_per_loop == 1:
return steps_per_loop
remainder_in_epoch = current_step % steps_per_epoch
if remainder_in_epoch != 0:
return min(steps_per_epoch - remainder_in_epoch, steps_per_loop)
else:
return steps_per_loop
def get_loss_fn(num_classes, loss_factor=1.0):
"""Gets the classification loss function."""
def classification_loss_fn(labels, logits):
"""Classification loss."""
labels = tf.squeeze(labels)
log_probs = tf.nn.log_softmax(logits, axis=-1)
one_hot_labels = tf.one_hot(
tf.cast(labels, dtype=tf.int32), depth=num_classes, dtype=tf.float32)
per_example_loss = -tf.reduce_sum(tf.cast(one_hot_labels, dtype=tf.float32) * log_probs, axis=-1)
loss = tf.reduce_mean(per_example_loss)
loss *= loss_factor
return loss
return classification_loss_fn
def run_customized_training_loop(
strategy,
model_fn,
loss_fn,
train_input_fn,
steps_per_loop,
steps_per_epoch,
epochs,
eval_input_fn,
eval_steps,
metric_fn,
):
total_training_steps = steps_per_epoch * epochs
train_input_data = train_input_fn()
train_iterator = iter(strategy.experimental_distribute_dataset(train_input_data))
with strategy.scope():
model = model_fn()
optimizer = model.optimizer
train_loss_metric = tf.keras.metrics.Mean(
'training_loss', dtype=tf.float32
)
eval_metrics = [metric_fn()]
train_metrics = [
metric.__class__.from_config(metric.get_config())
for metric in eval_metrics
]
# Collects training variables.
training_vars = model.trainable_variables
def _replicated_step(inputs):
"""Replicated training step."""
inputs, labels = inputs
with tf.GradientTape() as tape:
model_outputs = model(inputs, training=True)[0]
loss = loss_fn(labels, model_outputs)
grads = tape.gradient(loss, training_vars)
optimizer.apply_gradients(zip(grads, training_vars))
# For reporting, the metric takes the mean of losses.
train_loss_metric.update_state(loss)
for metric in train_metrics:
metric.update_state(labels, model_outputs)
@tf.function
def train_steps(iterator, steps):
for _ in tf.range(steps):
strategy.experimental_run_v2(_replicated_step, args=(next(iterator),))
def train_single_step(iterator):
strategy.experimental_run_v2(_replicated_step, args=(next(iterator),))
def test_step(iterator):
def _test_step_fn(inputs):
inputs, labels = inputs
model_outputs = model(inputs, training=False)
for metric in eval_metrics:
metric.update_state(labels, model_outputs)
strategy.experimental_run_v2(_test_step_fn, args=(next(iterator),))
train_single_step = tf.function(train_single_step)
test_step = tf.function(test_step)
def _run_evaluation(current_training_step, test_iterator):
for _ in range(eval_steps):
test_step(test_iterator)
current_step = optimizer.iterations.numpy()
while current_step < total_training_steps:
train_loss_metric.reset_states()
for metric in train_metrics + model.metrics:
metric.reset_states()
steps = _steps_to_run(current_step, steps_per_epoch, steps_per_loop)
if steps == 1:
train_single_step(train_iterator)
else:
train_steps(train_iterator, tf.convert_to_tensor(steps, dtype=tf.int32))
current_step += steps
train_loss = train_loss_metric.result().numpy().astype(float)
training_status = 'Train Step: %d/%d / loss = %s' % (current_step, total_training_steps, train_loss)
print(training_status)
return model
def run_customized_training(
tokenizer,
strategy,
dataset_path,
max_sequence_length,
train_batch_size,
eval_batch_size,
num_classes=10,
num_replicas=8,
steps_per_loop=10,
steps_per_epoch=10,
epochs=10,
eval_steps=10
):
train_input_fn = functools.partial(
create_dataset,
tokenizer,
dataset_path,
max_sequence_length,
train_batch_size
)
eval_input_fn = functools.partial(
create_dataset,
tokenizer,
dataset_path,
max_sequence_length,
eval_batch_size,
evaluate=True
)
def model_fn():
config = BertConfig.from_pretrained("bert-base-cased")
config.num_labels = 3
model = TFBertForSequenceClassification.from_pretrained("bert-base-cased", config=config)
optimizer = tf.keras.optimizers.Adam()
model.optimizer = optimizer
return model
loss_fn = get_loss_fn(num_classes, loss_factor=1.0/num_replicas)
def metric_fn():
return tf.keras.metrics.SparseCategoricalAccuracy('test_accuracy', dtype=tf.float32)
return run_customized_training_loop(
strategy=strategy,
model_fn=model_fn,
loss_fn=loss_fn,
train_input_fn=train_input_fn,
steps_per_loop=steps_per_loop,
steps_per_epoch=steps_per_epoch,
epochs=epochs,
eval_input_fn=eval_input_fn,
eval_steps=eval_steps,
metric_fn=metric_fn,
)
if __name__ == "__main__":
strategy, num_replicas = get_tpu()
tokenizer = BertTokenizer.from_pretrained("bert-base-cased")
input_meta_data_path = "gs://huggingface-bucket/transformers/outputs/MNLI_meta_data"
dataset_path = "/home/lysandre/transformers/examples/TPU/glue_data/MNLI"
with tf.io.gfile.GFile(input_meta_data_path, 'rb') as reader:
input_meta_data = json.loads(reader.read().decode('utf-8'))
max_sequence_length = input_meta_data["max_seq_length"]
num_classes = input_meta_data['num_labels']
epochs = 3
train_batch_size = 32
eval_batch_size = 32
train_data_size = input_meta_data["train_data_size"]
steps_per_epoch = int(train_data_size / train_batch_size)
steps_per_loop = 200
warmup_steps = int(epochs * train_data_size * 0.1 / train_batch_size)
eval_steps = int(math.ceil(input_meta_data['eval_data_size'] / eval_batch_size))
trained_model = run_customized_training(
tokenizer,
strategy,
dataset_path,
max_sequence_length,
train_batch_size,
eval_batch_size,
num_classes=num_classes,
num_replicas=num_replicas,
steps_per_loop=steps_per_loop,
steps_per_epoch=steps_per_epoch,
epochs=epochs,
eval_steps=eval_steps
)
@@ -0,0 +1,63 @@
# coding=utf-8
# Copyright 2018 The Open AI Team Authors and The HuggingFace Inc. team.
#
# Licensed under the Apache License, Version 2.0 (the "License");
# you may not use this file except in compliance with the License.
# You may obtain a copy of the License at
#
# http://www.apache.org/licenses/LICENSE-2.0
#
# Unless required by applicable law or agreed to in writing, software
# distributed under the License is distributed on an "AS IS" BASIS,
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
# See the License for the specific language governing permissions and
# limitations under the License.
""" === Under active development === Script to fine-tune GLUE on a TPU using keras' fit method
"""
from tpu_utils import get_tpu
import tensorflow as tf
from transformers import TFBertForSequenceClassification, BertTokenizer, glue_convert_examples_to_features
from time import time
import tensorflow_datasets
print("TF version: {}".format(tf.__version__))
num_epochs = 3
max_seq_length = 128
# The number of replicas should be obtained from the get_tpu() method, but the dataset pre-processing crashes if the
# TPU is loaded beforehand
num_replicas = 8
tokenizer = BertTokenizer.from_pretrained("bert-base-cased")
data = tensorflow_datasets.load('glue/mrpc')
train_dataset = glue_convert_examples_to_features(data['train'], tokenizer, max_seq_length, 'mrpc')
valid_dataset = glue_convert_examples_to_features(data['validation'], tokenizer, max_seq_length, 'mrpc')
total_train_batch_size = 32
train_batch_size_per_replica = total_train_batch_size / num_replicas
train_dataset = train_dataset.batch(total_train_batch_size)
assert train_batch_size_per_replica.is_integer()
total_valid_batch_size = 64
valid_batch_size_per_replica = total_valid_batch_size / num_replicas
valid_dataset = valid_dataset.batch(total_valid_batch_size)
assert valid_batch_size_per_replica.is_integer()
print('Fetched & created dataset.')
tpu, num_replicas = get_tpu()
with tpu.scope():
# Prepare training: Compile tf.keras model with optimizer, loss and learning rate schedule
optimizer = tf.keras.optimizers.Adam(learning_rate=3e-5, epsilon=1e-08, clipnorm=1.0)
loss = tf.keras.losses.SparseCategoricalCrossentropy(from_logits=True)
metric = tf.keras.metrics.SparseCategoricalAccuracy('accuracy')
model = TFBertForSequenceClassification.from_pretrained('bert-base-cased')
model.compile(optimizer=optimizer, loss=loss, metrics=[metric])
history = model.fit(train_dataset, epochs=2, steps_per_epoch=115,
validation_data=valid_dataset, validation_steps=7)
final_stats = model.evaluate(valid_dataset, steps=1)
print("Validation accuracy: ", final_stats[1])
+61
View File
@@ -0,0 +1,61 @@
# coding=utf-8
# Copyright 2018 The Open AI Team Authors and The HuggingFace Inc. team.
#
# Licensed under the Apache License, Version 2.0 (the "License");
# you may not use this file except in compliance with the License.
# You may obtain a copy of the License at
#
# http://www.apache.org/licenses/LICENSE-2.0
#
# Unless required by applicable law or agreed to in writing, software
# distributed under the License is distributed on an "AS IS" BASIS,
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
# See the License for the specific language governing permissions and
# limitations under the License.
""" === Under active development === Dataset load
"""
import tensorflow as tf
from transformers import glue_convert_examples_to_features, glue_processors
def create_dataset(tokenizer,
file_path,
seq_length,
batch_size,
is_training=True,
drop_remainder=False
):
processor = glue_processors["mnli"]()
examples = processor.get_dev_examples(file_path)
features = glue_convert_examples_to_features(examples, tokenizer, seq_length, 'mnli')
all_input_ids = tf.constant([f.input_ids for f in features])
all_attention_masks = tf.constant([f.attention_mask for f in features])
all_token_type_ids = tf.constant([f.token_type_ids for f in features])
all_labels = tf.constant([f.label for f in features])
dataset = tf.data.Dataset.from_tensor_slices(({
"input_ids": all_input_ids,
"attention_mask": all_attention_masks,
"token_type_ids": all_token_type_ids
}, all_labels))
dataset = dataset.batch(batch_size, drop_remainder=drop_remainder)
dataset = dataset.prefetch(1024)
return dataset
if __name__ == "__main__":
from transformers import BertTokenizer
train_data_path = "/home/lysandre/transformers/examples/TPU/glue_data/MNLI"
tokenizer = BertTokenizer.from_pretrained("bert-base-cased")
seq_length = 128
batch_size = 32
create_dataset(
tokenizer, train_data_path, seq_length, batch_size
)
+49
View File
@@ -0,0 +1,49 @@
# coding=utf-8
# Copyright 2018 The Open AI Team Authors and The HuggingFace Inc. team.
#
# Licensed under the Apache License, Version 2.0 (the "License");
# you may not use this file except in compliance with the License.
# You may obtain a copy of the License at
#
# http://www.apache.org/licenses/LICENSE-2.0
#
# Unless required by applicable law or agreed to in writing, software
# distributed under the License is distributed on an "AS IS" BASIS,
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
# See the License for the specific language governing permissions and
# limitations under the License.
""" === Under active development === Loading a TPUStrategy
Especially https://github.com/GoogleCloudPlatform/training-data-analyst/blob/tf2/courses/fast-and-lean-data-science/01_MNIST_TPU_Keras.ipynb
"""
import tensorflow as tf
def get_tpu():
tpu = None
try:
tpu = tf.distribute.cluster_resolver.TPUClusterResolver() # TPU detection
except ValueError as e:
print(e)
try:
tpu = tf.distribute.cluster_resolver.TPUClusterResolver(tpu="grpc://192.168.32.2:8470")
except ValueError as e:
print(e)
# Select appropriate distribution strategy
if tpu:
# TF 2.0 change here: experimental_connect_to_cluster and initialize_tpu_system are now necessary
tf.config.experimental_connect_to_cluster(tpu)
tf.tpu.experimental.initialize_tpu_system(tpu)
# TF 2.0 change here: steps_per_run does not exist anymore and is not needed
strategy = tf.distribute.experimental.TPUStrategy(tpu)
print('Running on TPU ', tpu.cluster_spec().as_dict()['worker'])
else:
strategy = tf.distribute.get_strategy() # default strategy that works on CPU and single GPU
print('Running on CPU or GPU')
print("Number of accelerators: ", strategy.num_replicas_in_sync)
return strategy, strategy.num_replicas_in_sync
+67 -40
View File
@@ -22,6 +22,7 @@ import glob
import logging
import os
import random
from time import time
import numpy as np
import torch
@@ -130,59 +131,72 @@ def train(args, train_dataset, model, tokenizer):
for _ in train_iterator:
epoch_iterator = tqdm(train_dataloader, desc="Iteration", disable=args.local_rank not in [-1, 0])
for step, batch in enumerate(epoch_iterator):
start = time()
model.train()
batch = tuple(t.to(args.device) for t in batch)
inputs = {'input_ids': batch[0],
'attention_mask': batch[1],
'labels': batch[3]}
print(batch[0].device)
if args.model_type != 'distilbert':
inputs['token_type_ids'] = batch[2] if args.model_type in ['bert', 'xlnet'] else None # XLM, DistilBERT and RoBERTa don't use segment_ids
intermediate = time()
print("Model took " + str(intermediate - start) + " to put on device.")
outputs = model(**inputs)
loss = outputs[0] # model outputs are always tuple in transformers (see doc)
if args.n_gpu > 1:
loss = loss.mean() # mean() to average on multi-gpu parallel training
if args.gradient_accumulation_steps > 1:
loss = loss / args.gradient_accumulation_steps
if args.fp16:
with amp.scale_loss(loss, optimizer) as scaled_loss:
scaled_loss.backward()
torch.nn.utils.clip_grad_norm_(amp.master_params(optimizer), args.max_grad_norm)
else:
loss.backward()
torch.nn.utils.clip_grad_norm_(model.parameters(), args.max_grad_norm)
# if args.n_gpu > 1:
# loss = loss.mean() # mean() to average on multi-gpu parallel training
# if args.gradient_accumulation_steps > 1:
# loss = loss / args.gradient_accumulation_steps
#
# if args.fp16:
# with amp.scale_loss(loss, optimizer) as scaled_loss:
# scaled_loss.backward()
# torch.nn.utils.clip_grad_norm_(amp.master_params(optimizer), args.max_grad_norm)
# else:
loss.backward()
# torch.nn.utils.clip_grad_norm_(model.parameters(), args.max_grad_norm)
tr_loss += loss.item()
if (step + 1) % args.gradient_accumulation_steps == 0:
optimizer.step()
scheduler.step() # Update learning rate schedule
model.zero_grad()
global_step += 1
# if (step + 1) % args.gradient_accumulation_steps == 0:
# optimizer.step()
# scheduler.step() # Update learning rate schedule
# model.zero_grad()
# global_step += 1
#
# if args.local_rank in [-1, 0] and args.logging_steps > 0 and global_step % args.logging_steps == 0:
# # Log metrics
# if args.local_rank == -1 and args.evaluate_during_training: # Only evaluate when single GPU otherwise metrics may not average well
# results = evaluate(args, model, tokenizer)
# for key, value in results.items():
# tb_writer.add_scalar('eval_{}'.format(key), value, global_step)
# tb_writer.add_scalar('lr', scheduler.get_lr()[0], global_step)
# tb_writer.add_scalar('loss', (tr_loss - logging_loss)/args.logging_steps, global_step)
# logging_loss = tr_loss
#
# if args.local_rank in [-1, 0] and args.save_steps > 0 and global_step % args.save_steps == 0:
# # Save model checkpoint
# output_dir = os.path.join(args.output_dir, 'checkpoint-{}'.format(global_step))
# if not os.path.exists(output_dir):
# os.makedirs(output_dir)
# model_to_save = model.module if hasattr(model, 'module') else model # Take care of distributed/parallel training
# model_to_save.save_pretrained(output_dir)
# torch.save(args, os.path.join(output_dir, 'training_args.bin'))
# logger.info("Saving model checkpoint to %s", output_dir)
#
# if args.max_steps > 0 and global_step > args.max_steps:
# epoch_iterator.close()
# break
if args.local_rank in [-1, 0] and args.logging_steps > 0 and global_step % args.logging_steps == 0:
# Log metrics
if args.local_rank == -1 and args.evaluate_during_training: # Only evaluate when single GPU otherwise metrics may not average well
results = evaluate(args, model, tokenizer)
for key, value in results.items():
tb_writer.add_scalar('eval_{}'.format(key), value, global_step)
tb_writer.add_scalar('lr', scheduler.get_lr()[0], global_step)
tb_writer.add_scalar('loss', (tr_loss - logging_loss)/args.logging_steps, global_step)
logging_loss = tr_loss
if args.local_rank in [-1, 0] and args.save_steps > 0 and global_step % args.save_steps == 0:
# Save model checkpoint
output_dir = os.path.join(args.output_dir, 'checkpoint-{}'.format(global_step))
if not os.path.exists(output_dir):
os.makedirs(output_dir)
model_to_save = model.module if hasattr(model, 'module') else model # Take care of distributed/parallel training
model_to_save.save_pretrained(output_dir)
torch.save(args, os.path.join(output_dir, 'training_args.bin'))
logger.info("Saving model checkpoint to %s", output_dir)
if args.max_steps > 0 and global_step > args.max_steps:
epoch_iterator.close()
break
end = time()
print("Model took " + str(end - start) + " for its forward pass.")
if args.max_steps > 0 and global_step > args.max_steps:
train_iterator.close()
break
@@ -379,6 +393,8 @@ def main():
parser.add_argument('--seed', type=int, default=42,
help="random seed for initialization")
parser.add_argument('--tpu', action='store_true',
help="Whether to use try and connect to a tpu")
parser.add_argument('--fp16', action='store_true',
help="Whether to use 16-bit (mixed) precision (through NVIDIA apex) instead of 32-bit")
parser.add_argument('--fp16_opt_level', type=str, default='O1',
@@ -393,6 +409,7 @@ def main():
if os.path.exists(args.output_dir) and os.listdir(args.output_dir) and args.do_train and not args.overwrite_output_dir:
raise ValueError("Output directory ({}) already exists and is not empty. Use --overwrite_output_dir to overcome.".format(args.output_dir))
# Setup distant debugging if needed
if args.server_ip and args.server_port:
# Distant debugging - see https://code.visualstudio.com/docs/python/debugging#_attach-to-a-local-script
@@ -410,6 +427,16 @@ def main():
device = torch.device("cuda", args.local_rank)
torch.distributed.init_process_group(backend='nccl')
args.n_gpu = 1
if args.tpu:
os.environ["TPU_IP_ADDRESS"] = "192.168.0.2"
os.environ["TPU_NAME"] = "node-1"
os.environ["XRT_TPU_CONFIG"] = "tpu_worker;0;192.168.0.2:8470"
import torch_xla
import torch_xla.core.xla_model as xm
device = xm.xla_device()
args.device = device
# Setup logging
+1
View File
@@ -256,6 +256,7 @@ def convert_examples_to_features(examples, tokenizer, max_seq_length,
start_offset += min(length, doc_stride)
for (doc_span_index, doc_span) in enumerate(doc_spans):
tokens_ = tokenizer.encode(query_tokens, )
tokens = []
token_to_orig_map = {}
token_is_max_context = {}
+1
View File
@@ -27,6 +27,7 @@ logger = logging.getLogger(__name__) # pylint: disable=invalid-name
try:
import tensorflow as tf
assert hasattr(tf, "__version__")
assert int(tf.__version__[0]) >= 2
_tf_available = True # pylint: disable=invalid-name
logger.info("TensorFlow version {} available.".format(tf.__version__))
+12 -2
View File
@@ -494,13 +494,23 @@ class TFBertMainLayer(tf.keras.layers.Layer):
position_ids = inputs.get('position_ids', position_ids)
head_mask = inputs.get('head_mask', head_mask)
assert len(inputs) <= 5, "Too many inputs."
if input_ids is None:
input_ids = inputs.get('input_word_ids')
token_type_ids = inputs.get('input_type_ids')
attention_mask = inputs.get('input_mask')
else:
input_ids = inputs
# TPUs have sharded objects
if attention_mask is None:
attention_mask = tf.fill(tf.shape(input_ids), 1)
attention_mask = tf.fill(tf.shape(input_ids.primary if hasattr(input_ids, 'primary') else input_ids), 1)
if token_type_ids is None:
token_type_ids = tf.fill(tf.shape(input_ids), 0)
token_type_ids = tf.fill(tf.shape(input_ids.primary if hasattr(input_ids, 'primary') else input_ids), 0)
# We create a 3D attention mask from a 2D tensor mask.
# Sizes are [batch_size, 1, 1, to_seq_length]