Skip to content

Instantly share code, notes, and snippets.

@Mageswaran1989
Created August 28, 2019 15:16
Show Gist options
  • Select an option

  • Save Mageswaran1989/facc3fc2a003807d029a914c721629db to your computer and use it in GitHub Desktop.

Select an option

Save Mageswaran1989/facc3fc2a003807d029a914c721629db to your computer and use it in GitHub Desktop.
Tensorflow 2.0 : Scalable Modleling with Estimator and Dataset
import glob
import os
import shutil
import tensorflow as tf
from tqdm import tqdm
import numpy as np
import matplotlib.pyplot as plt
import psutil
from absl import logging
import gc
from sklearn.datasets import make_regression
# import tensorflow.compat.v1 as tf
# tf.disable_v2_behavior()
logging.set_verbosity(logging.INFO)
"""
1. Create TFRecords
2. Define Model as part of Estimator
3. Read TFRecords into Dataset
4. Run Estimator with dataset
5. Collect memory stats
"""
TRAIN_DATA = out_dir = os.getcwd() + "/train_data"
VAL_DATA = out_dir = os.getcwd() + "/val_data"
MODEL_DIR = os.getcwd() + "/" + "nnet"
EXPORT_DIR = MODEL_DIR + "/" + "export"
NUM_SAMPLES_PER_FILE = 100000
NUM_FEATURES = 250
BATCH_SIZE = 128
def _int64_feature(value):
return tf.train.Feature(int64_list=tf.train.Int64List(value=[value]))
def _bytes_feature(value):
return tf.train.Feature(bytes_list=tf.train.BytesList(value=[value]))
def _float_feature(value):
return tf.train.Feature(float_list=tf.train.FloatList(value=[value]))
def _mat_feature(mat):
return tf.train.Feature(float_list=tf.train.FloatList(value=mat.flatten()))
def get_numpy_array_size(arr):
"""
Utility finction to get Numpy Array size
:param arr: Numpy Array
:return:
"""
size = (arr.size * arr.itemsize)/1024/1024
print("%d MBytes " % size)
return size
def _get_features(data, label):
"""
Converts numpy array as TF features
:param data: Numpy Array
:param label: Numpy Array
:return:
"""
return {
"data": _mat_feature(data),
"label": _float_feature(label)
}
def generate_tf_records(number_files, out_dir):
"""
Generates random data for Linear Regression and stores them as TFRecords
:param number_files: Number of TF records
:param out_dir: Out directory path
:return:
"""
if not os.path.exists(out_dir):
os.makedirs(out_dir)
for i in tqdm(range(number_files)):
file_path_name = os.path.join(out_dir, str(i) + ".tfrecords") # ~ 106MB
print(f"Writing to {file_path_name}")
if os.path.exists(file_path_name):
print(f"Found : {file_path_name}")
else:
# generate regression dataset
X, Y = make_regression(n_samples=NUM_SAMPLES_PER_FILE, n_features=NUM_FEATURES, noise=0.1)
# plot regression dataset
with tf.io.TFRecordWriter(file_path_name) as writer:
for x, y in zip(X, Y):
# (25000000 * 8) / 1024 /1024 ~ 190MB
# (100000 * 8) / 1024 /1024 ~ 0.76 MB
# get_numpy_array_size(features)
# get_numpy_array_size(label)
#create TF Features
features = tf.train.Features(feature=_get_features(data=x, label=y))
# create TF Example
example = tf.train.Example(features=features)
# print(example)
writer.write(example.SerializeToString())
def decode(serialized_example):
# define a parser
features = tf.io.parse_single_example(
serialized_example,
# Defaults are not specified since both keys are required.
features={
'data': tf.io.FixedLenFeature([1 * 250], tf.float32),
'label': tf.io.FixedLenFeature([1], tf.float32),
})
data = tf.reshape(
tf.cast(features['data'], tf.float32), shape=[1, 250])
label = tf.reshape(
tf.cast(features['label'], tf.float32), shape=[1])
return {"data": data}, label
def _get_dataset(data_path):
"""
Reads TFRecords, decode and batches them
:return: dataset
"""
_num_cores = 4
_batch_size = BATCH_SIZE
path = os.path.join(data_path, "*.tfrecords")
path = path.replace("//", "/")
files = tf.data.Dataset.list_files(path)
# files = glob.glob(pathname=path)
# TF dataset APIs
dataset = files.interleave(
tf.data.TFRecordDataset,
cycle_length=_num_cores,
num_parallel_calls=tf.data.experimental.AUTOTUNE)
# dataset = tf.data.TFRecordDataset(files, num_parallel_reads=_num_cores)
# dataset = dataset.shuffle(_batch_size*10, 42)
# Map the generator output as features as a dict and label
dataset = dataset.map(map_func=decode, num_parallel_calls=tf.data.experimental.AUTOTUNE)
dataset = dataset.batch(batch_size=_batch_size, drop_remainder=False)
dataset = dataset.prefetch(buffer_size=tf.data.experimental.AUTOTUNE)
iterator = dataset.make_one_shot_iterator()
batch_feats, batch_label = iterator.get_next()
return batch_feats, batch_label
def dataset_to_iterator(data_path):
"""
Reads the TFRecords and creates TF Datasets
:param data_path:
:return:
"""
_num_cores = 4
_batch_size = 128
path = os.path.join(data_path, "*.tfrecords")
path = path.replace("//", "/")
files = tf.data.Dataset.list_files(path)
# files = glob.glob(pathname=path)
# dataset = tf.data.TFRecordDataset(files, num_parallel_reads=_num_cores)
# TF dataset APIs
dataset = files.interleave(
tf.data.TFRecordDataset,
cycle_length=_num_cores,
num_parallel_calls=tf.data.experimental.AUTOTUNE)
dataset = dataset.shuffle(_batch_size*10, 42)
# Map the generator output as features as a dict and label
dataset = dataset.map(map_func=decode, num_parallel_calls=tf.data.experimental.AUTOTUNE)
dataset = dataset.batch(batch_size=_batch_size, drop_remainder=False)
dataset = dataset.prefetch(buffer_size=tf.data.experimental.AUTOTUNE)
# Create an iterator
iterator = dataset.make_one_shot_iterator()
# Create your tf representation of the iterator
image, label = iterator.get_next()
return image, label
class NNet(object):
def __init__(self):
pass
def __call__(self, features, labels, params, mode, config=None):
"""
Used for the :tf_main:`model_fn <estimator/Estimator#__init__>`
argument when constructing
:tf_main:`tf.estimator.Estimator <estimator/Estimator>`.
"""
return self._build(features, labels, params, mode, config=config)
def _get_optimizer(self, loss):
with tf.name_scope("optimizer") as scope:
global_step = tf.compat.v1.train.get_global_step()
learning_rate = tf.compat.v1.train.exponential_decay(0.001,
global_step,
decay_steps=100,
decay_rate=0.94,
staircase=True)
optimizer = tf.keras.optimizers.Adam(learning_rate=learning_rate,
beta_1=0.9,
beta_2=0.999,
epsilon=1e-7,
amsgrad=False,
name='Adam')
optimizer.iterations = tf.compat.v1.train.get_or_create_global_step()
# Get both the unconditional updates (the None part)
# and the input-conditional updates (the features part).
# update_ops = model.get_updates_for(None) + model.get_updates_for(features)
# Compute the minimize_op.
minimize_op = optimizer.get_updates(
loss,
tf.compat.v1.trainable_variables())[0]
train_op = tf.group(minimize_op)
return train_op
def _build(self, features, label, params, mode, config=None):
features = features['data']
net = tf.keras.layers.Dense(1024, activation='relu')(features)
net = tf.keras.layers.Dense(512, activation='relu')(net)
net = tf.keras.layers.Dense(256, activation='relu')(net)
net = tf.keras.layers.Dense(128, activation='relu')(net)
net = tf.keras.layers.Dense(64, activation='relu')(net)
net = tf.keras.layers.Dense(32, activation='relu')(net)
logits = tf.keras.layers.Dense(2, activation='softmax')(net)
classes = tf.math.greater(logits, 0.5)
loss = None
optimizer = None
predictions = {"probability" : logits, "classes" : classes}
if mode != tf.estimator.ModeKeys.PREDICT:
mse = tf.keras.losses.MeanSquaredError()
loss = mse(logits, label)
tf.summary.scalar('total_loss', loss)
optimizer = self._get_optimizer(loss=loss)
return tf.estimator.EstimatorSpec(
mode=mode,
predictions=predictions,
export_outputs={'predict': tf.estimator.export.PredictOutput(predictions)},
loss=loss,
train_op=optimizer,
eval_metric_ops=None)
def _init_tf_config(clear_model_data=False,
save_checkpoints_steps=NUM_SAMPLES_PER_FILE/10,
keep_checkpoint_max=5,
save_summary_steps=NUM_SAMPLES_PER_FILE/100,
log_step_count_steps=NUM_SAMPLES_PER_FILE/100):
run_config = tf.compat.v1.ConfigProto()
run_config.gpu_options.allow_growth = True
# run_config.gpu_options.per_process_gpu_memory_fraction = 0.50
run_config.allow_soft_placement = True
run_config.log_device_placement = False
model_dir = MODEL_DIR
if clear_model_data:
if os.path.exists(model_dir):
shutil.rmtree(model_dir)
_run_config = tf.estimator.RunConfig(session_config=run_config,
save_checkpoints_steps=save_checkpoints_steps,
keep_checkpoint_max=keep_checkpoint_max,
save_summary_steps=save_summary_steps,
model_dir=model_dir,
log_step_count_steps=log_step_count_steps)
return _run_config
def _get_train_spec(max_steps=None):
# Estimators expect an input_fn to take no arguments.
# To work around this restriction, we use lambda to capture the arguments and provide the expected interface.
return tf.estimator.TrainSpec(
input_fn=lambda: _get_dataset(data_path=TRAIN_DATA),
max_steps=max_steps,
hooks=None)
def _get_eval_spec(steps):
return tf.estimator.EvalSpec(
input_fn=lambda: _get_dataset(data_path=VAL_DATA),
steps=steps,
hooks=None)
def train(estimator, max_steps=None):
train_spec = _get_train_spec(max_steps=max_steps)
estimator.train(
input_fn=train_spec.input_fn,
hooks=train_spec.hooks,
max_steps=train_spec.max_steps)
def evaluate(estimator, steps=None, checkpoint_path=None):
eval_spec = _get_eval_spec(steps=steps)
estimator.evaluate(
input_fn=eval_spec.input_fn,
steps=eval_spec.steps,
hooks=eval_spec.hooks,
checkpoint_path=checkpoint_path)
def serving_input_receiver_fn():
inputs = {
"data": tf.compat.v1.placeholder(tf.float32, [None, 1, 250]),
}
return tf.estimator.export.ServingInputReceiver(inputs, inputs)
def export_model(estimator, model_export_path):
logging.info("Saving model to =======> {}".format(model_export_path))
if not os.path.exists(model_export_path):
os.makedirs(model_export_path)
estimator.export_saved_model(
model_export_path,
serving_input_receiver_fn=serving_input_receiver_fn)
def main():
memory_used = []
process = psutil.Process(os.getpid())
generate_tf_records(number_files=5, out_dir=TRAIN_DATA)
generate_tf_records(number_files=2, out_dir=VAL_DATA)
# print(dataset_to_iterator(data_path=TRAIN_DATA))
model = NNet()
estimator = tf.estimator.Estimator(model_fn=model, config=_init_tf_config(), params=None)
for epoch in tqdm(range(25)):
memory_used.append(process.memory_info()[0] / float(2 ** 20))
train(estimator=estimator)
evaluate(estimator=estimator)
plt.plot(memory_used)
plt.title('Evolution of memory')
plt.xlabel('iteration')
plt.ylabel('memory used (GB)')
plt.savefig("memory_usage.png")
plt.show()
export_model(estimator=estimator, model_export_path=EXPORT_DIR)
if __name__ == '__main__':
main()
"""
References:
- https://medium.com/mostly-ai/tensorflow-records-what-they-are-and-how-to-use-them-c46bc4bbb564
"""
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment