diff --git a/mnist_keras.py b/mnist_keras.py index 9fbbb70..aa42420 100644 --- a/mnist_keras.py +++ b/mnist_keras.py @@ -17,18 +17,18 @@ hvd.init() # Horovod: pin GPU to be used to process local rank (one GPU per process) -config = tf.ConfigProto() +config = tf.compat.v1.ConfigProto() config.gpu_options.allow_growth = True config.gpu_options.visible_device_list = str(hvd.local_rank()) -K.set_session(tf.Session(config=config)) +K.set_session(tf.compat.v1.Session(config=config)) -batch_size = 128 +batch_size = int(os.getenv('BATCH_SIZE', 128)) num_classes = 10 model_dir = os.path.abspath(os.environ.get('PS_MODEL_PATH', os.getcwd() + '/models') + '/horovod-mnist') export_dir = os.path.abspath(os.environ.get('PS_MODEL_PATH', os.getcwd() + '/models')) # Horovod: adjust number of epochs based on number of GPUs. -epochs = int(math.ceil(12.0 / hvd.size())) +epochs = int(math.ceil(int(os.getenv('TRAIN_EPOCHS', 12.0)) / hvd.size())) # Input image dimensions img_rows, img_cols = 28, 28 @@ -115,9 +115,12 @@ export_path = os.path.join(export_dir, str(version)) print('export_path = {}\n'.format(export_path)) - tf.saved_model.simple_save( - keras.backend.get_session(), - export_path, - inputs={'image': model.input}, - outputs={t.name:t for t in model.outputs}) +# tf.saved_model.simple_save( +# keras.backend.get_session(), +# export_path, +# inputs={'image': model.input}, +# outputs={t.name:t for t in model.outputs}) + + model.save(export_path, save_format='tf') + print('\nSaved model:') diff --git a/mnist_keras_tf2.py b/mnist_keras_tf2.py new file mode 100644 index 0000000..2c30c40 --- /dev/null +++ b/mnist_keras_tf2.py @@ -0,0 +1,98 @@ +# Copyright 2019 Uber Technologies, Inc. All Rights Reserved. +# +# 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. +# ============================================================================== + +import tensorflow as tf +import horovod.tensorflow.keras as hvd + +import os + +# Horovod: initialize Horovod. +hvd.init() + +# Horovod: pin GPU to be used to process local rank (one GPU per process) +gpus = tf.config.experimental.list_physical_devices('GPU') +for gpu in gpus: + tf.config.experimental.set_memory_growth(gpu, True) +if gpus: + tf.config.experimental.set_visible_devices(gpus[hvd.local_rank()], 'GPU') + +(mnist_images, mnist_labels), _ = \ + tf.keras.datasets.mnist.load_data(path='mnist-%d.npz' % hvd.rank()) + +dataset = tf.data.Dataset.from_tensor_slices( + (tf.cast(mnist_images[..., tf.newaxis] / 255.0, tf.float32), + tf.cast(mnist_labels, tf.int64)) +) +dataset = dataset.repeat().shuffle(10000).batch(int(os.getenv('BATCH_SIZE', 128))) + +export_dir = os.path.abspath(os.environ.get('PS_MODEL_PATH', os.getcwd() + '/models')) + +mnist_model = tf.keras.Sequential([ + tf.keras.layers.Conv2D(32, [3, 3], activation='relu'), + tf.keras.layers.Conv2D(64, [3, 3], activation='relu'), + tf.keras.layers.MaxPooling2D(pool_size=(2, 2)), + tf.keras.layers.Dropout(0.25), + tf.keras.layers.Flatten(), + tf.keras.layers.Dense(128, activation='relu'), + tf.keras.layers.Dropout(0.5), + tf.keras.layers.Dense(10, activation='softmax') +]) + +# Horovod: adjust learning rate based on number of GPUs. +opt = tf.optimizers.Adam(0.001 * hvd.size()) + +# Horovod: add Horovod DistributedOptimizer. +opt = hvd.DistributedOptimizer(opt) + +# Horovod: Specify `experimental_run_tf_function=False` to ensure TensorFlow +# uses hvd.DistributedOptimizer() to compute gradients. +mnist_model.compile(loss=tf.losses.SparseCategoricalCrossentropy(), + optimizer=opt, + metrics=['accuracy'], + experimental_run_tf_function=False) + +callbacks = [ + # Horovod: broadcast initial variable states from rank 0 to all other processes. + # This is necessary to ensure consistent initialization of all workers when + # training is started with random weights or restored from a checkpoint. + hvd.callbacks.BroadcastGlobalVariablesCallback(0), + + # Horovod: average metrics among workers at the end of every epoch. + # + # Note: This callback must be in the list before the ReduceLROnPlateau, + # TensorBoard or other metrics-based callbacks. + hvd.callbacks.MetricAverageCallback(), + + # Horovod: using `lr = 1.0 * hvd.size()` from the very beginning leads to worse final + # accuracy. Scale the learning rate `lr = 1.0` ---> `lr = 1.0 * hvd.size()` during + # the first three epochs. See https://arxiv.org/abs/1706.02677 for details. + hvd.callbacks.LearningRateWarmupCallback(warmup_epochs=3, verbose=1), +] + +# Horovod: save checkpoints only on worker 0 to prevent other workers from corrupting them. +if hvd.rank() == 0: + save_model_callback = tf.keras.callbacks.LambdaCallback( + on_train_end=lambda logs: [ + mnist_model.save(export_dir, save_format='tf')]) + + callbacks.append(tf.keras.callbacks.ModelCheckpoint(export_dir + '/checkpoint-{epoch}.h5')) + callbacks.append(save_model_callback) + +# Horovod: write logs on worker 0. +verbose = 1 if hvd.rank() == 0 else 0 + +# Train the model. +# Horovod: adjust number of steps based on number of GPUs. +mnist_model.fit(dataset, steps_per_epoch=int(os.getenv('STEPS_PER_EPOCH', 500)) // hvd.size(), callbacks=callbacks, epochs=int(os.getenv('TRAIN_EPOCHS', 24)), verbose=verbose)