Commit 2077f062 authored by Bharath Ramsundar's avatar Bharath Ramsundar
Browse files

First commit of transformer code.

parent ea66ca2b
Loading
Loading
Loading
Loading
+0 −15
Original line number Diff line number Diff line
@@ -19,7 +19,6 @@ from deepchem.utils.save import load_from_disk
from deepchem.utils.save import load_pandas_from_disk
from deepchem.utils import ScaffoldGenerator
from deepchem.featurizers.nnscore import NNScoreComplexFeaturizer
from pathos.multiprocessing import ProcessingPool
import multiprocessing as mp
from functools import partial
import multiprocess
@@ -220,11 +219,6 @@ class DataFeaturizer(object):
      features = []
      for ligand_protein_pdb_tuple in zip(ligand_pdbs, protein_pdbs):
        features.append(featurize_wrapper(ligand_protein_pdb_tuple))
    else:
      if worker_pool is None:
        worker_pool = ProcessingPool(mp.cpu_count())
        features = worker_pool.map(featurize_wrapper, 
                                   zip(ligand_pdbs, protein_pdbs))
    else:
      features = worker_pool.map_sync(featurize_wrapper, 
                                      zip(ligand_pdbs, protein_pdbs))
@@ -256,15 +250,6 @@ class DataFeaturizer(object):
        feature = featurizer.featurize([mol])
        return feature

      if worker_pool is None:
        dilled_featurizer = dill.dumps(featurizer)
        worker_pool = ProcessingPool(mp.cpu_count())
        featurize_wrapper_partial = partial(featurize_wrapper,
                                            dilled_featurizer=dilled_featurizer)
        features = []
        for smiles in sample_smiles:
          features.append(featurize_wrapper_partial(smiles))
      else:
      features = worker_pool.map_sync(featurize_wrapper, 
                                      sample_smiles)

+0 −443
Original line number Diff line number Diff line
"""
Process an input dataset into a format suitable for machine learning.
"""
from __future__ import print_function
from __future__ import division
from __future__ import unicode_literals
import os
import gzip
import pandas as pd
import numpy as np
import csv
from rdkit import Chem
from deepchem.featurizers.fingerprints import CircularFingerprint
from deepchem.featurizers.basic import RDKitDescriptors
from deepchem.utils.save import log
from deepchem.utils.save import save_to_disk
from deepchem.utils.save import load_from_disk
from deepchem.utils.save import load_pandas_from_disk
from deepchem.utils import ScaffoldGenerator
from deepchem.featurizers.nnscore import NNScoreComplexFeaturizer
import multiprocessing as mp
from pathos.multiprocessing import ProcessingPool
from functools import partial


def generate_scaffold(smiles, include_chirality=False):
  """Compute the Bemis-Murcko scaffold for a SMILES string."""
  mol = Chem.MolFromSmiles(smiles)
  engine = ScaffoldGenerator(include_chirality=include_chirality)
  scaffold = engine.get_scaffold(mol)
  return scaffold

def _check_validity(compounds_df):
  """Ensure that columns of compound_df contain required elements."""
  if not set(FeaturizedSamples.colnames).issubset(compounds_df.keys()):
    raise ValueError("Compound dataframe does not contain required columns")

def _process_field(val):
  """Parse data in a field."""
  if isinstance(val, float) or isinstance(val, np.ndarray):
    return val
  elif isinstance(val, list):
    return [_process_field(elt) for elt in val]
  elif isinstance(val, str):
    try:
      return float(val)
    except ValueError:
      return val
  else:
    raise ValueError("Field of unrecognized type: %s" % str(val))

def _get_input_type(input_file):
  """Get type of input file. Must be csv/pkl.gz/sdf file."""
  filename, file_extension = os.path.splitext(input_file)
  # If gzipped, need to compute extension again
  if file_extension == ".gz":
    filename, file_extension = os.path.splitext(filename)
  if file_extension == ".csv":
    return "csv"
  elif file_extension == ".pkl":
    return "pandas-pickle"
  elif file_extension == ".joblib":
    return "pandas-joblib"
  else:
    raise ValueError("Unrecognized extension %s" % file_extension)

def _get_fields(input_file):
  """Get the names of fields and field_types for input data."""
  # If CSV input, assume that first row contains labels
  input_type = _get_input_type(input_file)
  if input_type == "csv":
    with open(input_file, "rb") as inp_file_obj:
      return csv.reader(inp_file_obj).next()
  elif input_type == "pandas-joblib":
    df = load_from_disk(input_file)
    return df.keys()
  elif input_type == "pandas-pickle":
    df = load_pickle_from_disk(input_file)
    return df.keys()
  else:
    raise ValueError("Unrecognized extension for %s" % input_file)

class DataFeaturizer(object):
  """
  Handles loading/featurizing of chemical samples (datapoints).

  Currently knows how to load csv-files/pandas-dataframes/SDF-files. Writes a
  dataframe object to disk as output.
  """

  def __init__(self, tasks, smiles_field, split_field=None,
               id_field=None, threshold=None, user_specified_features=None,
               protein_pdb_field=None, ligand_pdb_field=None,
               ligand_mol2_field=None, 
               compound_featurizers=[], complex_featurizers=[],
               verbose=False, log_every_n=1000):
    """Extracts data from input as Pandas data frame"""
    if not isinstance(tasks, list):
      raise ValueError("tasks must be a list.")
    self.tasks = tasks
    self.smiles_field = smiles_field
    self.split_field = split_field
    if id_field is None:
      self.id_field = smiles_field
    else:
      self.id_field = id_field
    self.threshold = threshold
    self.protein_pdb_field = protein_pdb_field
    self.ligand_pdb_field = ligand_pdb_field
    self.ligand_mol2_field = ligand_mol2_field
    self.user_specified_features = user_specified_features
    self.compound_featurizers = compound_featurizers
    self.complex_featurizers = complex_featurizers
    self.verbose = verbose
    self.log_every_n = log_every_n

  def featurize(self, input_file, feature_dir, samples_dir, shard_size=128):
    """Featurize provided file and write to specified location."""
    input_type = _get_input_type(input_file)

    log("Loading raw samples now.", self.verbose)
    raw_df = load_pandas_from_disk(input_file)
    fields = raw_df.keys()
    log("Loaded raw data frame from file.", self.verbose)

    def process_raw_sample_helper(row, fields, input_type):
      return self._process_raw_sample(input_type, row, fields)
    process_raw_sample_helper_partial = partial(process_raw_sample_helper,
                                                fields=fields,
                                                input_type=input_type)


    #processed_rows = raw_df.apply(process_raw_sample_helper_partial, axis=1)
    raw_df = raw_df.apply(process_raw_sample_helper_partial, axis=1, reduce=False)
    #raw_df = pd.DataFrame.from_records(processed_rows)

    nb_sample = raw_df.shape[0]
    interval_points = np.linspace(
        0, nb_sample, np.ceil(float(nb_sample)/shard_size)+1, dtype=int)
    shard_files = []
    for j in range(len(interval_points)-1):
      log("Sharding and standardizing into shard-%s / %s shards" % (str(j+1), len(interval_points)-1), self.verbose)
      raw_df_shard = raw_df.iloc[range(interval_points[j], interval_points[j+1])]
      
      df = self._standardize_df(raw_df_shard) 
      log("Aggregating User-Specified Features", self.verbose)
      self._add_user_specified_features(df)

      for compound_featurizer in self.compound_featurizers:
        log("Currently feauturizing feature_type: %s"
            % compound_featurizer.__class__.__name__, self.verbose)
        self._featurize_compounds(df, compound_featurizer)

      for complex_featurizer in self.complex_featurizers:
        log("Currently feauturizing feature_type: %s"
            % complex_featurizer.__class__.__name__, self.verbose)
        self._featurize_complexes(df, complex_featurizer)

      shard_out = os.path.join(feature_dir, "features_shard%d.joblib" % j)
      save_to_disk(df, shard_out)
      shard_files.append(shard_out)

    featurizers = self.compound_featurizers + self.complex_featurizers
    samples = FeaturizedSamples(samples_dir=samples_dir, featurizers=featurizers, 
                                dataset_files=shard_files,
                                reload_data=False)

    return samples

  def _process_raw_sample(self, input_type, row, fields):
    """Extract information from row data."""
    data = {}
    if input_type == "csv":
      for ind, field in enumerate(fields):
        data[field] = _process_field(row[ind])
    elif input_type in ["pandas-pickle", "pandas-joblib"]:
      for field in fields:
        data[field] = _process_field(row[field])
    else:
      raise ValueError("Unrecognized input_type")
    if self.threshold is not None:
      for task in self.tasks:
        raw = _process_field(data[task])
        if not isinstance(raw, float):
          raise ValueError("Cannot threshold non-float fields.")
        data[field] = 1 if raw > self.threshold else 0
    return data

  def _standardize_df(self, ori_df):
    """Copy specified columns to new df with standard column names."""
    df = pd.DataFrame(ori_df[[self.id_field]])
    df.columns = ["mol_id"]
    df["smiles"] = ori_df[[self.smiles_field]]
    for task in self.tasks:
      df[task] = ori_df[[task]]
    if self.split_field is not None:
      df["split"] = ori_df[[self.split_field]]
    if self.protein_pdb_field is not None:
      df["protein_pdb"] = ori_df[[self.protein_pdb_field]]
    if self.ligand_pdb_field is not None:
      df["ligand_pdb"] = ori_df[[self.ligand_pdb_field]]
    if self.ligand_mol2_field is not None:
      df["ligand_mol2"] = ori_df[[self.ligand_mol2_field]]
    return df

  def _featurize_complexes(self, df, featurizer, parallel=True):
    """Generates circular fingerprints for dataset."""
    protein_pdbs = list(df["protein_pdb"])
    ligand_pdbs = list(df["ligand_pdb"])
    complexes = zip(ligand_pdbs, protein_pdbs)

    def featurize_wrapper(ligand_protein_pdb_tuple):
      ligand_pdb, protein_pdb = ligand_protein_pdb_tuple
      molecule_features = featurizer.featurize_complexes([ligand_pdb], [protein_pdb])
      return featurizer.featurize_complexes([ligand_pdb], [protein_pdb])

    features = ProcessingPool(mp.cpu_count()).map(featurize_wrapper, 
                                                  zip(ligand_pdbs, protein_pdbs))
    #features = featurize_wrapper(zip(ligand_pdbs, protein_pdbs))
    df[featurizer.__class__.__name__] = list(features)

  def _featurize_compounds(self, df, featurizer, parallel=True):    
    """Featurize individual compounds.

       Given a featurizer that operates on individual chemical compounds 
       or macromolecules, compute & add features for that compound to the 
       features dataframe
    """
    sample_smiles = df["smiles"].tolist()

    if not parallel:
      features = []
      for ind, smiles in enumerate(sample_smiles):
        if ind % self.log_every_n == 0:
          log("Featurizing sample %d" % ind, self.verbose)
        mol = Chem.MolFromSmiles(smiles)
        features.append(featurizer.featurize([mol]))
    else:
      def featurize_wrapper(smiles):
        mol = Chem.MolFromSmiles(smiles)
        return featurizer.featurize([mol])

      features = ProcessingPool(mp.cpu_count()).map(featurize_wrapper, 
                                                    sample_smiles)

    df[featurizer.__class__.__name__] = features

  def _add_user_specified_features(self, df):
    """Merge user specified features. 

      Merge features included in dataset provided by user
      into final features dataframe
    """
    if self.user_specified_features is not None:
      log("Adding user-defined features.", self.verbose)
      features_data = []
      for _, row in df.iterrows():
        # pandas rows are tuples (row_num, row_data)
        feature_list = []
        for feature_name in self.user_specified_features:
          feature_list.append(row[feature_name])
        features_data.append({"user-specified-features": np.array(feature_list)})
      df["user-specified-features"] = pd.DataFrame(features_data)


def map_function(data_tuple, featurizer):
  featurizer = NNScoreComplexFeaturizer()
  ind, ligand_pdb, protein_pdb = data_tuple
  #print("Mapping on ind %d" % ind)
  #print("ind, type(ligand_pdb), type(protein_pdb): %s " % str((ind, type(ligand_pdb), type(protein_pdb))))
  return featurizer.featurize_complexes([ligand_pdb], [protein_pdb])

class FeaturizedSamples(object):
  """
  Wrapper class for featurized data on disk.
  """
  # The standard columns for featurized data.
  colnames = ["mol_id", "smiles", "split"]
  optional_colnames = ["protein_pdb", "ligand_pdb", "ligand_mol2"]

  def __init__(self, samples_dir, featurizers, dataset_files=None, 
               overwrite=True, reload_data=False):
    """
    Initialiize FeaturizedSamples

    If samples_dir does not exist, must specify dataset_files. Then samples_dir
    is created and populated. If samples_dir exists (created by previous call to
    FeaturizedSamples), then dataset_files cannot be specified. If overwrite is
    set and dataset_files is provided, will overwrite old dataset_files with
    new.
    """
    self.dataset_files = dataset_files
    self.feature_types = (
        ["user-specified-features"] + 
        [featurizer.__class__.__name__ for featurizer in featurizers])

    self.featurizers = featurizers

    if not os.path.exists(samples_dir):
      os.makedirs(samples_dir)
    self.samples_dir = samples_dir
    if os.path.exists(self._get_compounds_filename()) and reload_data:
      compounds_df = load_from_disk(self._get_compounds_filename())
    else:
      compounds_df = self._get_compounds()
      # compounds_df is not altered by any method after initialization, so it's
      # safe to keep a copy in memory and on disk.
      save_to_disk(compounds_df, self._get_compounds_filename())
    _check_validity(compounds_df)
    self.compounds_df = compounds_df

    if os.path.exists(self._get_dataset_paths_filename()):
      if dataset_files is not None:
        if overwrite:
          save_to_disk(dataset_files, self._get_dataset_paths_filename())
        else:
          raise ValueError("Can't change dataset_files already stored on disk")
      self.dataset_files = load_from_disk(self._get_dataset_paths_filename())
    else:
      save_to_disk(dataset_files, self._get_dataset_paths_filename())

  def _get_compounds_filename(self):
    """
    Get standard location for file listing compounds in this dataframe.
    """
    return os.path.join(self.samples_dir, "compounds.joblib")

  def _get_dataset_paths_filename(self):
    """
    Get standard location for file listing dataset_files.
    """
    return os.path.join(self.samples_dir, "datasets.joblib")

  def _get_compounds(self):
    """
    Create dataframe containing metadata about compounds.
    """
    compound_rows = []
    for dataset_file in self.dataset_files:
      df = load_from_disk(dataset_file)
      compound_ids = list(df["mol_id"])
      smiles = list(df["smiles"])
      if "split" in df.keys():
        splits = list(df["split"])
      else:
        splits = [None] * len(smiles)
      compound_rows += [list(elt) for elt in zip(compound_ids, smiles, splits)]

    compounds_df = pd.DataFrame(compound_rows,
                                columns=("mol_id", "smiles", "split"))
    return compounds_df

  def _set_compound_df(self, df):
    """Internal method used to replace compounds_df."""
    _check_validity(df)
    save_to_disk(df, self._get_compounds_filename())
    self.compounds_df = df

  # TODO(rbharath): Might this be inefficient?
  def itersamples(self):
    """
    Provides an iterator over samples.

    Each sample from the iterator is a dataframe of samples.
    """
    compound_ids = set(list(self.compounds_df["mol_id"]))
    for df_file in self.dataset_files:
      df = load_from_disk(df_file)
      visible_inds = []
      for ind, row in df.iterrows():
        if row["mol_id"] in compound_ids:
          visible_inds.append(ind)
      yield df.loc[visible_inds]

  def train_test_split(self, splittype, train_dir, test_dir, seed=None,
                       frac_train=.8):
    """
    Splits self into train/test sets and returns two FeaturizedDatsets
    """
    if splittype == "random":
      train_inds, test_inds = self._train_test_random_split(seed=seed, frac_train=frac_train)
    elif splittype == "scaffold":
      train_inds, test_inds = self._train_test_scaffold_split(frac_train=frac_train)
    elif splittype == "specified":
      train_inds, test_inds = self._train_test_specified_split()
    else:
      raise ValueError("improper splittype.")
    train_samples = FeaturizedSamples(samples_dir=train_dir, 
                                      dataset_files=self.dataset_files,
                                      featurizers=self.featurizers)
    train_samples._set_compound_df(self.compounds_df.iloc[train_inds])
    test_samples = FeaturizedSamples(samples_dir=test_dir, 
                                     dataset_files=self.dataset_files,
                                     featurizers=self.featurizers)
    test_samples._set_compound_df(self.compounds_df.iloc[test_inds])

    return train_samples, test_samples

  def _train_test_random_split(self, seed=None, frac_train=.8):
    """
    Splits internal compounds randomly into train/test.
    """
    np.random.seed(seed)
    train_cutoff = frac_train * len(self.compounds_df)
    shuffled = np.random.permutation(range(len(self.compounds_df)))
    return shuffled[:train_cutoff], shuffled[train_cutoff:]

  def _train_test_scaffold_split(self, frac_train=.8):
    """
    Splits internal compounds into train/test by scaffold.
    """
    scaffolds = {}
    for ind, row in self.compounds_df.iterrows():
      scaffold = generate_scaffold(row["smiles"])
      if scaffold not in scaffolds:
        scaffolds[scaffold] = [ind]
      else:
        scaffolds[scaffold].append(ind)
    # Sort from largest to smallest scaffold sets
    scaffold_sets = [scaffold_set for (scaffold, scaffold_set) in
                     sorted(scaffolds.items(), key=lambda x: -len(x[1]))]
    train_cutoff = frac_train * len(self.compounds_df)
    train_inds, test_inds = [], []
    for scaffold_set in scaffold_sets:
      if len(train_inds) + len(scaffold_set) > train_cutoff:
        test_inds += scaffold_set
      else:
        train_inds += scaffold_set
    return train_inds, test_inds

  def _train_test_specified_split(self):
    """
    Splits internal compounds into train/test by user-specification.
    """
    train_inds, test_inds = [], []
    for ind, row in self.compounds_df.iterrows():
      if row["split"].lower() == "train":
        train_inds.append(ind)
      elif row["split"].lower() == "test":
        test_inds.append(ind)
      else:
        raise ValueError("Missing required split information.")
    return train_inds, test_inds
+21 −104
Original line number Diff line number Diff line
@@ -49,8 +49,6 @@ class Dataset(object):
        if retval is not None:
          metadata_rows.append(retval)

      # TODO(rbharath): FeaturizedSamples should not be responsible for
      # X-transform, X_sums, etc. Move that stuff over to Dataset.
      self.metadata_df = pd.DataFrame(
          metadata_rows,
          columns=('df_file', 'task_names', 'ids',
@@ -58,20 +56,26 @@ class Dataset(object):
                   'w',
                   'X_sums', 'X_sum_squares', 'X_n',
                   'y_sums', 'y_sum_squares', 'y_n'))
      save_to_disk(
          self.metadata_df, self._get_metadata_filename())
      # input/output transforms not specified yet, so
      # self.transforms = (input_transforms, output_transforms) =>
      self.transforms = ([], [])
      save_to_disk(
          self.transforms, self._get_transforms_filename())
      self.save()
      #save_to_disk(
      #    self.metadata_df, self._get_metadata_filename())
      ## input/output transforms not specified yet, so
      ## self.transforms = (input_transforms, output_transforms) =>
      #self.transforms = ([], [])
      #save_to_disk(
      #    self.transforms, self._get_transforms_filename())
    else:
      if os.path.exists(self._get_metadata_filename()):
        self.metadata_df = load_from_disk(self._get_metadata_filename())
        self.transforms = load_from_disk(self._get_transforms_filename())
        #self.transforms = load_from_disk(self._get_transforms_filename())
      else:
        raise ValueError("No metadata found.")

  def save_to_disk():
    """Save dataset to disk."""
    save_to_disk(
        self.metadata_df, self._get_metadata_filename())

  def get_task_names(self):
    """
    Gets learning tasks associated with this dataset.
@@ -96,11 +100,11 @@ class Dataset(object):
    metadata_filename = os.path.join(self.data_dir, "metadata.joblib")
    return metadata_filename

  def _get_transforms_filename(self):
    """
    Get standard location for stored transforms.
    """
    return os.path.join(self.data_dir, "transforms.joblib")
  #def _get_transforms_filename(self):
  #  """
  #  Get standard location for stored transforms.
  #  """
  #  return os.path.join(self.data_dir, "transforms.joblib")

  def get_number_shards(self):
    """
@@ -122,45 +126,6 @@ class Dataset(object):
      ids = load_from_disk(row['ids'])
      yield (X, y, w, ids)

  def transform(self, input_transforms, output_transforms, parallel=False):
    """
    Transforms all internally stored data.

    Adds X-transform, y-transform columns to metadata.
    """
    (normalize_X, truncate_x, normalize_y, truncate_y, log_X, log_y) = (
        False, False, False, False, False, False)

    if "truncate" in input_transforms:
      truncate_x = True
    if "normalize" in input_transforms:
      normalize_X = True
    if "log" in input_transforms:
      log_X = True

    if "normalize" in output_transforms:
      normalize_y = True
    if "log" in output_transforms:
      log_y = True

    # Store input_transforms/output_transforms so the dataset remembers its state.

    X_means, X_stds, y_means, y_stds = self._transform(normalize_X, normalize_y,
                                                       truncate_x, truncate_y,
                                                       log_X, log_y,
                                                       parallel=parallel)
    nrow = self.metadata_df.shape[0]
    # TODO(rbharath): These lines are puzzling. Better way to avoid storage
    # duplication here?
    self.metadata_df['X_means'] = [X_means for _ in range(nrow)]
    self.metadata_df['X_stds'] = [X_stds for _ in range(nrow)]
    self.metadata_df['y_means'] = [y_means for _ in range(nrow)]
    self.metadata_df['y_stds'] = [y_stds for _ in range(nrow)]
    save_to_disk(
        self.metadata_df, self._get_metadata_filename())
    self.transforms = (input_transforms, output_transforms)
    save_to_disk(
        self.transforms, self._get_transforms_filename())

  def get_label_means(self):
    """Return pandas series of label means."""
@@ -180,60 +145,12 @@ class Dataset(object):
    (_, output_transforms) = self.transforms
    return output_transforms

  def _transform(self, normalize_X=True, normalize_y=True,
                 truncate_X=True, truncate_y=True,
                 log_X=False, log_y=False, parallel=False):
    """Helper to (parallel) transform all indexed data."""
  def compute_statistics(self):
    """Computes statistics of this dataset"""
    df = self.metadata_df
    trunc = 5.0
    X_means, X_stds, y_means, y_stds = compute_mean_and_std(df)
    indices = range(0, df.shape[0])
    transform_row_partial = partial(_transform_row, df=df, normalize_X=normalize_X,
                                    normalize_y=normalize_y, truncate_X=truncate_X,
                                    truncate_y=truncate_y, log_X=log_X,
                                    log_y=log_y, X_means=X_means, X_stds=X_stds,
                                    y_means=y_means, y_stds=y_stds, trunc=trunc)
    if parallel:
      pool = mp.Pool(int(mp.cpu_count()/4))
      pool.map(transform_row_partial, indices)
      pool.terminate()
    else:
      for index in indices:
        transform_row_partial(index)

    return X_means, X_stds, y_means, y_stds
    
def _transform_row(i, df, normalize_X, normalize_y, truncate_X, truncate_y,
                   log_X, log_y, X_means, X_stds, y_means, y_stds, trunc):
  """
  Transforms the data (X, y, w,...) in a single row.

  Writes X-transforme,d y-transformed to disk.
  """
  row = df.iloc[i]
  X = load_from_disk(row['X'])
  if normalize_X or log_X:
    if normalize_X:
      # Turns NaNs to zeros
      X = np.nan_to_num((X - X_means) / X_stds)
      if truncate_X:
        X[X > trunc] = trunc
        X[X < (-1.0*trunc)] = -1.0 * trunc
    if log_X:
      X = np.log(X)
  save_to_disk(X, row['X-transformed'])

  y = load_from_disk(row['y'])
  if normalize_y or log_y:
    if normalize_y:
      y = np.nan_to_num((y - y_means) / y_stds)
      if truncate_y:
        y[y > trunc] = trunc
        y[y < (-1.0*trunc)] = -1.0 * trunc
    if log_y:
      y = np.log(y)
  save_to_disk(y, row['y-transformed'])

def compute_sums_and_nb_sample(tensor, W=None):
  """
  Computes sums, squared sums of tensor along axis 0.
+30 −24

File changed.

Preview size limit exceeded, changes collapsed.