Commit 8d98ffb4 authored by Bharath Ramsundar's avatar Bharath Ramsundar
Browse files

Able to run tf models in prediction mode. Still crashes due to tf/deepchem eval mismatch

parent d412d029
Loading
Loading
Loading
Loading
+58 −80
Original line number Diff line number Diff line
@@ -52,8 +52,6 @@ class TensorflowModel(Model):
    add_output_ops
    build
    Eval
    read_input / _read_input_generator
    BatchInputGenerator (if you want to use mr_eval/EvalBatch)
    training_cost 

  Subclasses must set the following attributes:
@@ -93,8 +91,8 @@ class TensorflowModel(Model):
  def __init__(self,
               task_types,
               model_params,
               train,
               logdir,
               train=False,
               logdir=None,
               graph=None,
               summary_writer=None):
    self.model_params = model_params 
@@ -117,12 +115,16 @@ class TensorflowModel(Model):
    self._name_scopes = {}

    with self.graph.as_default():
      model_ops.set_training(train)
      self.placeholder_root = 'placeholders'
      with tf.name_scope(self.placeholder_root) as scope:
        self.placeholder_scope = scope
        self.valid = tf.placeholder(tf.bool,
                                    shape=[model_params["batch_size"]],
                                    name='valid')
      print("TensorflowModel.__init__")
      print("self.valid")
      print(self.valid)

    if "num_classification_tasks" in model_params:
      num_classification_tasks = model_params["num_classification_tasks"]
@@ -192,6 +194,9 @@ class TensorflowModel(Model):
            tf.placeholder(tf.float32, shape=[self.model_params["batch_size"]],
                           name='weights_%d' % task)))
    self.weights = weights
    print("TensorflowClassiifer.add_weight_placeholders")
    print("self.weights")
    print(self.weights)

  def add_labels_and_weights(self):
    """Add Placeholders for labels and weights.
@@ -209,41 +214,6 @@ class TensorflowModel(Model):
    self.add_label_placeholders()
    self.add_weight_placeholders()

  def read_input(self, input_pattern):
    """Read input data and return a generator for minibatches.

    Args:
      input_pattern: Input file pattern.

    Returns:
      A generator that yields a dict for feeding a single batch to Placeholders
      in the graph.

    Raises:
      NotImplementedError: if not overridden by concrete subclass.
    """
    raise NotImplementedError('Must be overridden by concrete subclass')

  def _read_input_generator(self, names, tensors):
    """Generator that constructs feed_dict for minibatches.

    read_input cannot be a generator because any reading ops will not be added to
    the graph until .next() is called (which is too late). Instead, read_input
    should perform any necessary graph construction and then return this
    generator.

    Args:
      names: A list of tensor names.
      tensors: A list of tensors to evaluate.

    Yields:
      A dict for feeding a single batch to Placeholders in the graph.

    Raises:
      NotImplementedError: if not overridden by concrete subclass.
    """
    raise NotImplementedError('Must be overridden by concrete subclass')

  def _shared_name_scope(self, name):
    """Returns a singleton TensorFlow scope with the given name.

@@ -255,6 +225,7 @@ class TensorflowModel(Model):
      tf.name_scope with the provided name.
    """
    if name not in self._name_scopes:
      with self.graph.as_default():
        with tf.name_scope(name) as scope:
          self._name_scopes[name] = scope

@@ -277,7 +248,8 @@ class TensorflowModel(Model):
    raise NotImplementedError('Must be overridden by concrete subclass')

  def training_cost(self):
    self.RequireAttributes(['output', 'labels', 'weights'])
    with self.graph.as_default():
      self.require_attributes(['output', 'labels', 'weights'])
      epsilon = 1e-3  # small float to avoid dividing by zero
      model_params = self.model_params
      weighted_costs = []  # weighted costs for each example
@@ -340,13 +312,15 @@ class TensorflowModel(Model):

  def setup(self):
    """Add ops common to training/eval to the graph."""
    with self.graph.as_default():
      with tf.name_scope('core_model'):
        self.build()
      self.add_labels_and_weights()
      self.global_step = tf.Variable(0, name='global_step', trainable=False)

  def MergeUpdates(self):
  def merge_updates(self):
    """Group updates into a single op."""
    with self.graph.as_default():
      updates = tf.get_default_graph().get_collection('updates')
      if updates:
        self.updates = tf.group(*updates, name='updates')
@@ -371,6 +345,7 @@ class TensorflowModel(Model):
    Returns:
    A summary op.
    """
    with self.graph.as_default():
      return tf.merge_all_summaries()


@@ -397,20 +372,20 @@ class TensorflowModel(Model):
    Raises:
      AssertionError: If model is not in training mode.
    """
    model_ops.SetTraining(True)
    #assert model_ops.IsTraining()
    with self.graph.as_default():
      assert model_ops.is_training()
      self.setup()
      self.training_cost()
    self.MergeUpdates()
    self.RequireAttributes(['loss', 'global_step', 'updates'])
      self.merge_updates()
      self.require_attributes(['loss', 'global_step', 'updates'])
      if summaries:
      self.AddSummaries()
        self.add_summaries()
      train_op = self.get_training_op()
      summary_op = self.get_summary_op()
      no_op = tf.no_op()
      # TODO(rbharath): This should probably be uncommented!
    #tf.train.write_graph(
    #    tf.get_default_graph().as_graph_def(), self.logdir, 'train.pbtxt')
      tf.train.write_graph(
          tf.get_default_graph().as_graph_def(), self.logdir, 'train.pbtxt')
      self.summary_writer.add_graph(tf.get_default_graph().as_graph_def())
      last_checkpoint_time = time.time()
      last_summary_time = time.time()
@@ -477,13 +452,14 @@ class TensorflowModel(Model):
    Args:
      checkpoint: string. Path to checkpoint file.
    """
    with self.graph.as_default():
      if self._restored_model:
        return

    with self.graph.as_default():
      self.setup()
      self.add_output_ops()  # add softmax heads
      saver = tf.train.Saver(tf.variables.all_variables())
      #saver = tf.train.Saver(tf.variables.all_variables())
      saver = tf.train.Saver()
      saver.restore(self._get_shared_session(),
                    tf_utils.ParseCheckpoint(checkpoint))
      self.global_step_number = int(self._get_shared_session().run(self.global_step))
@@ -494,6 +470,8 @@ class TensorflowModel(Model):
    """
    Loads model from disk. Thin wrapper around restore() for consistency.
    """
    with self.graph.as_default():
      assert not model_ops.is_training()
      if model_dir != self.logdir:
        raise ValueError("Cannot load from directory that is not logdir.")
      last_checkpoint = self._find_last_checkpoint()
@@ -505,13 +483,15 @@ class TensorflowModel(Model):
    for filename in os.listdir(self.logdir):
      # checkpoints look like logdir/model.ckpt-N
      # self._save_path is "logdir/model.ckpt"
      if self._save_path in filename:
      if os.path.basename(self._save_path) in filename:
        try:
          N = int(filename.split("-")[-1])
          if N > highest_num:
            highest_num = N
            last_checkpoint = filename
    return last_checkpoint
          
        except ValueError:
          pass
    return os.path.join(self.logdir, last_checkpoint)
          

  def Eval(self, input_generator, checkpoint, metrics=None):
@@ -530,7 +510,7 @@ class TensorflowModel(Model):
        values for each task.
    """
    self.Restore(checkpoint)
    output, labels, weights = self.ModelOutput(input_generator)
    output, labels, weights = self.get_model_output(input_generator)
    y_true, y_pred = self.ParseModelOutput(output, labels, weights)

    # keep counts for each class as a sanity check
@@ -574,29 +554,13 @@ class TensorflowModel(Model):
      computed_metrics.append(metric_value)
    return computed_metrics

  def EvalBatch(self, input_batch):
    """Runs inference on the provided batch of input.

    Args:
      input_batch: iterator of input with len self.model_params.batch_size.

    Returns:
      Tuple of three numpy arrays with shape num_examples x num_tasks (x ...):
        output: Model predictions.
        labels: True labels.
        weights: Example weights.
    """
    with self.graph.as_default():
      return self.ModelOutput(self.BatchInputGenerator(input_batch))

  def ModelOutput(self, input_generator):
  def predict(self, dataset, transformers):
    """Return model output for the provided input.

    Restore(checkpoint) must have previously been called on this object.

    Args:
      input_generator: Generator that returns a feed_dict for feeding
        Placeholders in the model graph.
      dataset: deepchem.datasets.dataset object.

    Returns:
      Tuple of three numpy arrays with shape num_examples x num_tasks (x ...):
@@ -610,9 +574,10 @@ class TensorflowModel(Model):
      AssertionError: If model is not in evaluation mode.
      ValueError: If output and labels are not both 3D or both 2D.
    """
    assert not model_ops.IsTraining()
    with self.graph.as_default():
      assert not model_ops.is_training()
      assert self._restored_model
    self.RequireAttributes(['output', 'labels', 'weights'])
      self.require_attributes(['output', 'labels', 'weights'])

      # run eval data through the model
      num_tasks = self.num_tasks
@@ -622,7 +587,9 @@ class TensorflowModel(Model):
        batches_per_summary = 1000
        seconds_per_summary = 0
        batch_count = -1.0
      for feed_dict in input_generator:
        #for feed_dict in input_generator:
        for (X_b, y_b, w_b, ids_b) in dataset.iterbatches(self.model_params["batch_size"]):
          feed_dict = self.construct_feed_dict(X_b, y_b, w_b, ids_b)
          batch_start = time.time()
          batch_count += 1
          data = self._get_shared_session().run(
@@ -645,6 +612,8 @@ class TensorflowModel(Model):
                (batch_output.shape, batch_labels.shape))
          batch_weights = batch_weights.transpose((1, 0))
          valid = feed_dict[self.valid.name]
          print("valid")
          print(valid)
          # only take valid outputs
          if np.count_nonzero(~valid):
            batch_output = batch_output[valid]
@@ -656,7 +625,7 @@ class TensorflowModel(Model):

          # Writes summary for tracking eval progress.
          seconds_per_summary += (time.time() - batch_start)
        self.RequireAttributes(['summary_writer'])
          self.require_attributes(['summary_writer'])
          if batch_count % batches_per_summary == 0:
            mean_seconds_per_batch = seconds_per_summary / batches_per_summary
            seconds_per_summary = 0
@@ -687,6 +656,7 @@ class TensorflowModel(Model):
      name: String name for this group of metrics. Useful for organizing
        metrics calculated using different subsets of the data.
    """
    with self.graph.as_default():
      # create a DataFrame to hold results
      data = dict()
      if counts is not None:
@@ -705,7 +675,7 @@ class TensorflowModel(Model):
        pickle.dump(df, f, pickle.HIGHEST_PROTOCOL)

      # write a summary for each metric
    self.RequireAttributes(['summary_writer'])
      self.require_attributes(['summary_writer'])
      with tf.Session(self.master):
        summaries = []
        prefix = '' if name is None else '%s - ' % name
@@ -723,7 +693,7 @@ class TensorflowModel(Model):
                                        global_step=global_step)
        self.summary_writer.flush()

  def RequireAttributes(self, attrs):
  def require_attributes(self, attrs):
    """Require class attributes to be defined.

    Args:
@@ -737,8 +707,9 @@ class TensorflowModel(Model):
        raise AssertionError(
            'self.%s must be defined by a concrete subclass' % attr)

  def AddSummaries(self):
  def add_summaries(self):
    """Add summaries for model parameters."""
    with self.graph.as_default():
      for var in tf.trainable_variables():
        if 'BatchNormalize' in var.name:
          continue
@@ -759,7 +730,7 @@ class TensorflowClassifier(TensorflowModel):
  default_metrics = ['auc']

  def get_task_type(self):
    return "classifier"
    return "classification"

  def cost(self, logits, labels, weights):
    """Calculate single-task training cost for a batch of examples.
@@ -783,6 +754,7 @@ class TensorflowClassifier(TensorflowModel):
    Returns:
      A list of tensors with shape batch_size containing costs for each task.
    """
    with self.graph.as_default():
      weighted_costs = super(TensorflowClassifier, self).training_cost()  # calculate loss
      epsilon = 1e-3  # small float to avoid dividing by zero
      model_params = self.model_params
@@ -834,6 +806,7 @@ class TensorflowClassifier(TensorflowModel):

  def add_output_ops(self):
    """Replace logits with softmax outputs."""
    with self.graph.as_default():
      softmax = []
      with tf.name_scope('inference'):
        for i, logits in enumerate(self.output):
@@ -866,6 +839,7 @@ class TensorflowClassifier(TensorflowModel):
    Placeholders are wrapped in identity ops to avoid the error caused by
    feeding and fetching the same tensor.
    """
    with self.graph.as_default():
      model_params = self.model_params
      batch_size = model_params["batch_size"]
      num_classes = model_params["num_classes"]
@@ -876,6 +850,9 @@ class TensorflowClassifier(TensorflowModel):
              tf.placeholder(tf.float32, shape=[batch_size, num_classes],
                             name='labels_%d' % task)))
      self.labels = labels
      print("TensorflowClassiifer.add_label_placeholders")
      print("self.labels")
      print(self.labels)

  def ParseModelOutput(self, output, labels, weights):
    """Parse model output to get true and predicted values for each task.
@@ -955,6 +932,7 @@ class TensorflowRegressor(TensorflowModel):
    Placeholders are wrapped in identity ops to avoid the error caused by
    feeding and fetching the same tensor.
    """
    with self.graph.as_default():
      batch_size = self.model_params["batch_size"]
      labels = []
      for task in xrange(self.num_tasks):
+14 −4
Original line number Diff line number Diff line
@@ -91,6 +91,7 @@ class TensorflowMultiTaskClassifier(TensorflowClassifier):
      mol_features: Molecule descriptor (e.g. fingerprint) tensor with shape
        batch_size x num_features.
    """
    with self.graph.as_default():
      with tf.name_scope(self.placeholder_scope):
        self.mol_features = tf.placeholder(
            tf.float32,
@@ -156,11 +157,20 @@ class TensorflowMultiTaskClassifier(TensorflowClassifier):
      w_b: np.ndarray of shape (batch_size, num_tasks)
      ids_b: List of length (batch_size) with datapoint identifiers.
    """ 
    feed_dict = {}
    feed_dict[self.mol_features] = X_b
    orig_dict = {}
    orig_dict["mol_features"] = X_b
    for task in xrange(self.num_tasks):
      feed_dict[self.labels[task]] = to_one_hot(y_b[:, task])
      feed_dict[self.weights[task]] = w_b[:, task]
      orig_dict["labels_%d" % task] = to_one_hot(y_b[:, task])
      orig_dict["weights_%d" % task] = w_b[:, task]
    orig_dict["valid"] = np.ones((self.model_params["batch_size"],), dtype=bool)
    return self._get_feed_dict(orig_dict)

  # TODO(rbharath): This explicit manipulation of scopes is ugly. Is there a
  # better design here?
  def _get_feed_dict(self, named_values):
    feed_dict = {}
    for name, value in named_values.iteritems():
      feed_dict['{}/{}:0'.format(self.placeholder_root, name)] = value
    return feed_dict

  def ReadInput(self, input_pattern, input_data_types=None):
+17 −8
Original line number Diff line number Diff line
@@ -14,6 +14,9 @@
# See the License for the specific language governing permissions and
# limitations under the License.
"""Ops for graph construction."""
from __future__ import print_function
from __future__ import division
from __future__ import unicode_literals


import tensorflow as tf
@@ -24,6 +27,8 @@ from tensorflow.python.platform import gfile
from tensorflow.python.platform import logging

from deepchem.models.tensorflow_models import utils as model_utils
import sys
import traceback


def AddBias(tensor, init=None, name=None):
@@ -97,7 +102,7 @@ def BatchNormalize(tensor, convolution, mask=None, epsilon=0.001,
    # moving averages from training.
    mean_moving_average = MovingAverage(mean, global_step, decay)
    variance_moving_average = MovingAverage(variance, global_step, decay)
    if not IsTraining():
    if not is_training():
      mean = mean_moving_average
      variance = variance_moving_average

@@ -168,7 +173,7 @@ def Dropout(tensor, dropout_prob, training_only=True):
  if not dropout_prob:
    return tensor  # do nothing
  keep_prob = 1.0 - dropout_prob
  if IsTraining() or not training_only:
  if is_training() or not training_only:
    tensor = tf.nn.dropout(tensor, keep_prob)
  return tensor

@@ -204,8 +209,7 @@ def FullyConnectedLayer(tensor, size, weight_init=None, bias_init=None,
    b = tf.Variable(bias_init, name='b')
    return tf.nn.xw_plus_b(tensor, w, b)


def IsTraining():
def is_training():
  """Determine whether the default graph is in training mode.

  Returns:
@@ -215,20 +219,25 @@ def IsTraining():
    ValueError: If the 'train' collection in the default graph does not contain
      exactly one element.
  """
  train = tf.get_collection('train')
  #traceback.print_stack(file=sys.stdout) 
  train = tf.get_collection("train")
  print("is_training()")
  print("train")
  print(train)
  if not train:
    raise ValueError('Training mode is not set. Please call SetTraining.')
    raise ValueError('Training mode is not set. Please call set_training.')
  elif len(train) > 1:
    raise ValueError('Training mode has more than one setting.')
  return train[0]


def SetTraining(train):
def set_training(train):
  """Set the training mode of the default graph.

  This operation may only be called once for a given graph.

  Args:
    graph: Tensorflow graph. 
    train: If True, graph is in training mode.

  Raises:
@@ -236,7 +245,7 @@ def SetTraining(train):
  """
  if tf.get_collection('train'):
    raise AssertionError('Training mode already set: %s' %
                         tf.get_collection('train'))
                         graph.get_collection('train'))
  tf.add_to_collection('train', train)


+20 −8
Original line number Diff line number Diff line
@@ -42,22 +42,32 @@ class TestAPI(unittest.TestCase):
    self.samples_dir = tempfile.mkdtemp()
    self.train_dir = tempfile.mkdtemp()
    self.test_dir = tempfile.mkdtemp()
    self.model_dir = tempfile.mkdtemp()
    #self.model_dir = tempfile.mkdtemp()
    self.model_dir = "/home/rbharath/model_dir" 
    if not os.path.exists(self.model_dir):
      os.makedirs(self.model_dir)

  def tearDown(self):
    shutil.rmtree(self.feature_dir)
    shutil.rmtree(self.samples_dir)
    shutil.rmtree(self.train_dir)
    shutil.rmtree(self.test_dir)
    shutil.rmtree(self.model_dir)
    #shutil.rmtree(self.model_dir)

  def _create_model(self, train_dataset, test_dataset, model, transformers):
  def _create_model(self, train_dataset, test_dataset, model, transformers,
                    test_model_creator=None):
    """Helper method to create model for test."""

    # Fit model
    # Fit trained model
    model.fit(train_dataset)
    model.save(self.model_dir)

    # Now create test model
    if test_model_creator is not None:
      test_model = test_model_creator()
      test_model.load(self.model_dir)
      model = test_model

    # Eval model on train
    evaluator = Evaluator(model, train_dataset, transformers, verbose=True)
    with tempfile.NamedTemporaryFile() as train_csv_out:
@@ -383,8 +393,10 @@ class TestAPI(unittest.TestCase):
      "learning_rate": .001,
      "data_shape": train_dataset.get_data_shape()
    }
    train = True
    logdir = self.model_dir
    model = TensorflowMultiTaskClassifier(
        task_types, model_params, train, logdir)
    self._create_model(train_dataset, test_dataset, model, transformers)
        task_types, model_params, train=True, logdir=self.model_dir)
    def test_model_creator():
      return TensorflowMultiTaskClassifier(
          task_types, model_params, train=False, logdir=self.model_dir)
    self._create_model(train_dataset, test_dataset, model, transformers,
                       test_model_creator=test_model_creator)
+4 −6
Original line number Diff line number Diff line
@@ -92,18 +92,16 @@ class Evaluator(object):
    performance_df = pd.DataFrame(columns=colnames)

    for i, task_name in enumerate(self.task_names):
      print("task_name")
      print(task_name)
      print("pred_y_df.keys()")
      print(pred_y_df.keys())
      y = pred_y_df[task_name].values
      y_pred = pred_y_df["%s_pred" % task_name].values
      if threshold is not None:
        # TODO(rbharath): This is a hack. More structured approach?
        y = pred_y_df[task_name+"_raw"].values
        #print("y")
        #print(y)
        #print("y_pred_orig")
        #print(y_pred)
        y_pred = threshold_predictions(y_pred, threshold)
        #print("y_pred")
        #print(y_pred)
      w = pred_y_df["%s_weight" % task_name].values

      if task_type == "classification":