In [1]:
%matplotlib inline

from __future__ import absolute_import
from __future__ import division
from __future__ import print_function

import collections
import math
import os
import time
import random
import zipfile

import numpy as np
from six.moves import urllib
from six.moves import xrange  # pylint: disable=redefined-builtin
import tensorflow as tf
from matplotlib import pylab



## Step 1: Download the Data

In [2]:
url = 'http://mattmahoney.net/dc/'


def maybe_download(filename, expected_bytes):
    """Download a file if not present, and make sure it's the right size."""
    if not os.path.exists(filename):
        filename, _ = urllib.request.urlretrieve(url + filename, filename)
        statinfo = os.stat(filename)
        if statinfo.st_size == expected_bytes:
            print('Found and verified', filename)
        else:
            print(statinfo.st_size)
            raise Exception(
                'Failed to verify {}. \
                Can you get to it with a browser?'.format(
                    filename
                )
            )
    return filename

filename = maybe_download('text8.zip', 31344016)

### Read the data into a list of strings

In [3]:
def read_data(filename):
    """Extract the first file enclosed in a zip file as a list of words"""
    with zipfile.ZipFile(filename) as f:
        data = f.read(f.namelist()[0]).split()
    return data

words = read_data(filename)
print('Data size', len(words))

Data size 17005207


## Step 2: Build the dictionary and replace rare words with UNK token

In [4]:
vocabulary_size = 50000


def build_dataset(words):
    count = [['UNK', -1]]
    count.extend(collections.Counter(words).most_common(vocabulary_size - 1))
    dictionary = dict()
    for word, _ in count:
        dictionary[word] = len(dictionary)
    data = list()
    unk_count = 0
    for word in words:
        if word in dictionary:
            index = dictionary[word]
        else:
            index = 0  # dictionary['UNK']
            unk_count += 1
        data.append(index)
    count[0][1] = unk_count
    reverse_dictionary = dict(zip(dictionary.values(), dictionary.keys()))
    return data, count, dictionary, reverse_dictionary

data, count, dictionary, reverse_dictionary = build_dataset(words)
del words  # Hint to reduce memory.
print('Most common words (+UNK)', count[:5])
print('Sample data', data[:10], [reverse_dictionary[i] for i in data[:10]])

data_index = 0

Most common words (+UNK) [['UNK', 418391], ('the', 1061396), ('of', 593677), ('and', 416629), ('one', 411764)]
Sample data [5239, 3084, 12, 6, 195, 2, 3137, 46, 59, 156] ['anarchism', 'originated', 'as', 'a', 'term', 'of', 'abuse', 'first', 'used', 'against']


## Step 3: Function to generate a training batch for the skip-gram model

In [5]:
def generate_batch(batch_size, num_skips, skip_window):
    global data_index
    assert batch_size % num_skips == 0
    assert num_skips <= 2 * skip_window
    batch = np.ndarray(shape=(batch_size), dtype=np.int32)
    labels = np.ndarray(shape=(batch_size, 1), dtype=np.int32)
    span = 2 * skip_window + 1  # [ skip_window target skip_window ]
    buffer = collections.deque(maxlen=span)
    for _ in range(span):
        buffer.append(data[data_index])
        data_index = (data_index + 1) % len(data)
    for i in range(batch_size // num_skips):
        target = skip_window  # target label at the center of the buffer
        targets_to_avoid = [skip_window]
        for j in range(num_skips):
            while target in targets_to_avoid:
                target = random.randint(0, span - 1)
            targets_to_avoid.append(target)
            batch[i * num_skips + j] = buffer[skip_window]
            labels[i * num_skips + j, 0] = buffer[target]
        buffer.append(data[data_index])
        data_index = (data_index + 1) % len(data)
    return batch, labels

batch, labels = generate_batch(batch_size=8, num_skips=2, skip_window=1)
for i in range(8):
    print(
        batch[i],
        reverse_dictionary[batch[i]],
        '->',
        labels[i, 0],
        reverse_dictionary[labels[i, 0]]
    )

3084 originated -> 5239 anarchism
3084 originated -> 12 as
12 as -> 3084 originated
12 as -> 6 a
6 a -> 12 as
6 a -> 195 term
195 term -> 6 a
195 term -> 2 of


## Step 4: Build and train a skip-gram model

### Fetch the cluster definition

In [6]:
import ast

cluster_config = ast.literal_eval(os.environ.get('CLUSTER_CONFIG'))
cluster_spec = tf.train.ClusterSpec(cluster_config)
workers = ['/job:worker/task:{}'.format(i) for i in range(len(cluster_config['worker']))]
param_servers = ['/job:ps/task:{}'.format(i) for i in range(len(cluster_config['ps']))]

print(cluster_config)

{'ps': ['ps-0.default.svc.cluster.local:8080', 'ps-1.default.svc.cluster.local:8080', 'ps-2.default.svc.cluster.local:8080', 'ps-3.default.svc.cluster.local:8080'], 'worker': ['worker-0.default.svc.cluster.local:8080', 'worker-1.default.svc.cluster.local:8080', 'worker-2.default.svc.cluster.local:8080', 'worker-3.default.svc.cluster.local:8080', 'worker-4.default.svc.cluster.local:8080', 'worker-5.default.svc.cluster.local:8080', 'worker-6.default.svc.cluster.local:8080', 'worker-7.default.svc.cluster.local:8080'], 'master': ['master-0.default.svc.cluster.local:8080']}


### Define Hyperparameters and validation set

In [7]:
batch_size = 128
embedding_size = 128  # Dimension of the embedding vector.
skip_window = 1       # How many words to consider left and right.
num_skips = 2         # How many times to reuse an input to generate a label.

# We pick a random validation set to sample nearest neighbors. Here we limit
# validation samples to the words that have a low numeric ID, which by
# construction are also the most frequent.
valid_size = 16     # Random set of words to evaluate similarity on.
valid_window = 100  # Only pick dev samples in the head of the distribution.
valid_examples = np.random.choice(valid_window, valid_size, replace=False)
num_sampled = 64    # Number of negative examples to sample.

### Build the model

In [41]:
graph = tf.Graph()
with graph.as_default():
    valid_dataset = tf.constant(valid_examples, dtype=tf.int32)

    # Input data.
    train_inputs = tf.placeholder(tf.int32, shape=[batch_size])
    train_labels = tf.placeholder(tf.int32, shape=[batch_size, 1])

    # Make one minibatch of data for each worker
    train_inputs_list = tf.split(0, len(workers), train_inputs)
    train_labels_list = tf.split(0, len(workers), train_labels)


    # Create a variable embedding for each parameter server
    embeddings = []
    for param_server in param_servers:
        with tf.device(param_server):
            embeddings.append(
                tf.Variable(tf.random_uniform(
                    [vocabulary_size, embedding_size],
                    -1.0,
                    1.0
                ))
            )

    with tf.device(tf.train.replica_device_setter(cluster=cluster_spec)):
        # Construct the variables for the NCE loss
        nce_weights = tf.Variable(
            tf.truncated_normal([vocabulary_size, embedding_size],
                                stddev=1.0 / math.sqrt(embedding_size)))
        nce_biases = tf.Variable(tf.zeros([vocabulary_size]))
        global_step = tf.Variable(0)

    losses = []
    summaries = []
    # Assign computational tasks to each worker
    for i, worker in enumerate(workers):
        with tf.device(worker):
            # Look up embeddings for inputs.
            embed = tf.nn.embedding_lookup(embeddings, train_inputs_list[i])
            # Compute the average NCE loss for the batch.
            # tf.nce_loss automatically draws a new sample of the negative
            # labels each time we evaluate the loss.
            loss = tf.reduce_mean(
                tf.nn.nce_loss(nce_weights, nce_biases, embed, train_labels_list[i],
                               num_sampled, vocabulary_size))
            losses.append(loss)
            summaries.append(tf.scalar_summary("loss-{}".format(i), loss))


    average_loss_op = tf.add_n(losses) / tf.convert_to_tensor(len(losses), dtype=tf.float32)
    summaries.append(tf.scalar_summary("average-loss", average_loss_op))
    
    summary_op = tf.merge_summary(summaries)
    init_op = tf.initialize_all_variables()
    saver = tf.train.Saver(tf.all_variables(), sharded=True)

    train_op = tf.train.GradientDescentOptimizer(1.0, use_locking=True).minimize(
            average_loss_op, global_step=global_step)

    # Compute the cosine similarity between minibatch examples
    # and the embeddings
    average_embeddings = tf.add_n(embeddings) / tf.convert_to_tensor(len(embeddings), dtype=tf.float32)

    norm = tf.sqrt(tf.reduce_sum(
               tf.square(average_embeddings), 1, keep_dims=True))
    normalized_embeddings = average_embeddings / norm
    valid_embeddings = tf.nn.embedding_lookup(average_embeddings, valid_dataset)
    similarity = tf.matmul(
        valid_embeddings, normalized_embeddings, transpose_b=True)

print("Graph Definition Complete")
        

Graph Definition Complete


## Step 5: Begin Training

In [53]:
num_steps = 100000

save_dir = '/var/log/checkpoints/word2vec/{}'.format(int(time.time()))
save_dir_prefix = save_dir + '/model'

if not os.path.exists(save_dir):
    os.makedirs(save_dir)

sm = tf.train.SessionManager(graph=graph)
last_report_time = time.time()
last_report_step = 0
average_loss_total = 0
step = 0

summary_writer = tf.train.SummaryWriter(
    '/var/log/tensorflow/word2vec/{}/'.format(int(time.time())),
    graph=graph
)

with sm.prepare_session('grpc://localhost:8080',
                        init_op=init_op,
                        saver=saver,
                        checkpoint_dir=save_dir) as session:
    while step < num_steps:
        try:
            batch_inputs, batch_labels = generate_batch(
                batch_size, num_skips, skip_window)
            feed_dict = {train_inputs: batch_inputs, train_labels: batch_labels}

            # We perform one update step by evaluating the optimizer op
            # Also evaluate the training summary op
            # Average the workers' losses, and advance the global step
            _, average_loss, step, summary = session.run(
                [train_op, average_loss_op, global_step, summary_op],
                feed_dict=feed_dict
            )

            summary_writer.add_summary(summary, global_step=step)
            average_loss_total += average_loss
            cur_time = time.time()

            if cur_time > last_report_time + 10:
                # The average loss is an estimate of the loss over the last 2000
                # batches.
                steps = step - last_report_step
                print("Average loss at step {}: {} \t Steps/Second: {}".format(
                        step,
                        average_loss_total / steps,
                        steps / (cur_time - last_report_time)
                ))
                last_report_time = cur_time
                last_report_step = step
                average_loss_total = 0
                saver.save(session, save_dir_prefix, global_step=step)
                
        except tf.errors.UnavailableError as e:
            print("You Killed a Worker")
            session = sm.prepare_session(
                'grpc://localhost:8080',
                init_op=init_op,
                saver=saver,
                checkpoint_dir=save_dir)
            
    print("   Nearest words:")
    sim = similarity.eval()
    for i in xrange(valid_size):
        valid_word = reverse_dictionary[valid_examples[i]]
        top_k = 8  # number of nearest neighbors
        nearest = (-sim[i, :]).argsort()[1:top_k+1]
        print("      {}: {}".format(
            valid_word, ', '.join(
                reverse_dictionary[nearest[k]] for k in range(top_k)
            ))
        )
    final_embeddings = normalized_embeddings.eval()


Average loss at step 67: 233.033552255 	 Steps/Second: 6.67815903507


NotFoundError: /var/log/checkpoints/word2vec/1464465701/model-67.tempstate14588595220314914492
	 [[Node: save/save = SaveSlices[T=[DT_FLOAT, DT_FLOAT, DT_FLOAT, DT_FLOAT, DT_FLOAT, DT_FLOAT, DT_INT32], _device="/job:master/replica:0/task:0/cpu:0"](_recv_save/Const_0, save/save/tensor_names, save/save/shapes_and_slices, Variable_S1165, Variable_1_S1167, Variable_2_S1169, Variable_3_S1171, Variable_4_S1173, Variable_5_S1175, Variable_6_S1177)]]
Caused by op u'save/save', defined at:
  File "/usr/lib/python2.7/runpy.py", line 162, in _run_module_as_main
    "__main__", fname, loader, pkg_name)
  File "/usr/lib/python2.7/runpy.py", line 72, in _run_code
    exec code in run_globals
  File "/usr/local/lib/python2.7/dist-packages/ipykernel/__main__.py", line 3, in <module>
    app.launch_new_instance()
  File "/usr/local/lib/python2.7/dist-packages/traitlets/config/application.py", line 596, in launch_instance
    app.start()
  File "/usr/local/lib/python2.7/dist-packages/ipykernel/kernelapp.py", line 442, in start
    ioloop.IOLoop.instance().start()
  File "/usr/local/lib/python2.7/dist-packages/zmq/eventloop/ioloop.py", line 162, in start
    super(ZMQIOLoop, self).start()
  File "/usr/local/lib/python2.7/dist-packages/tornado/ioloop.py", line 883, in start
    handler_func(fd_obj, events)
  File "/usr/local/lib/python2.7/dist-packages/tornado/stack_context.py", line 275, in null_wrapper
    return fn(*args, **kwargs)
  File "/usr/local/lib/python2.7/dist-packages/zmq/eventloop/zmqstream.py", line 440, in _handle_events
    self._handle_recv()
  File "/usr/local/lib/python2.7/dist-packages/zmq/eventloop/zmqstream.py", line 472, in _handle_recv
    self._run_callback(callback, msg)
  File "/usr/local/lib/python2.7/dist-packages/zmq/eventloop/zmqstream.py", line 414, in _run_callback
    callback(*args, **kwargs)
  File "/usr/local/lib/python2.7/dist-packages/tornado/stack_context.py", line 275, in null_wrapper
    return fn(*args, **kwargs)
  File "/usr/local/lib/python2.7/dist-packages/ipykernel/kernelbase.py", line 276, in dispatcher
    return self.dispatch_shell(stream, msg)
  File "/usr/local/lib/python2.7/dist-packages/ipykernel/kernelbase.py", line 228, in dispatch_shell
    handler(stream, idents, msg)
  File "/usr/local/lib/python2.7/dist-packages/ipykernel/kernelbase.py", line 391, in execute_request
    user_expressions, allow_stdin)
  File "/usr/local/lib/python2.7/dist-packages/ipykernel/ipkernel.py", line 199, in do_execute
    shell.run_cell(code, store_history=store_history, silent=silent)
  File "/usr/local/lib/python2.7/dist-packages/IPython/core/interactiveshell.py", line 2723, in run_cell
    interactivity=interactivity, compiler=compiler, result=result)
  File "/usr/local/lib/python2.7/dist-packages/IPython/core/interactiveshell.py", line 2825, in run_ast_nodes
    if self.run_code(code, result):
  File "/usr/local/lib/python2.7/dist-packages/IPython/core/interactiveshell.py", line 2885, in run_code
    exec(code_obj, self.user_global_ns, self.user_ns)
  File "<ipython-input-41-57ab53fec0cd>", line 56, in <module>
    saver = tf.train.Saver(tf.all_variables())
  File "/usr/local/lib/python2.7/dist-packages/tensorflow/python/training/saver.py", line 832, in __init__
    restore_sequentially=restore_sequentially)
  File "/usr/local/lib/python2.7/dist-packages/tensorflow/python/training/saver.py", line 500, in build
    save_tensor = self._AddSaveOps(filename_tensor, vars_to_save)
  File "/usr/local/lib/python2.7/dist-packages/tensorflow/python/training/saver.py", line 197, in _AddSaveOps
    save = self.save_op(filename_tensor, vars_to_save)
  File "/usr/local/lib/python2.7/dist-packages/tensorflow/python/training/saver.py", line 149, in save_op
    tensor_slices=[vs.slice_spec for vs in vars_to_save])
  File "/usr/local/lib/python2.7/dist-packages/tensorflow/python/ops/io_ops.py", line 172, in _save
    tensors, name=name)
  File "/usr/local/lib/python2.7/dist-packages/tensorflow/python/ops/gen_io_ops.py", line 341, in _save_slices
    name=name)
  File "/usr/local/lib/python2.7/dist-packages/tensorflow/python/ops/op_def_library.py", line 661, in apply_op
    op_def=op_def)
  File "/usr/local/lib/python2.7/dist-packages/tensorflow/python/framework/ops.py", line 2154, in create_op
    original_op=self._default_original_op, op_def=op_def)
  File "/usr/local/lib/python2.7/dist-packages/tensorflow/python/framework/ops.py", line 1154, in __init__
    self._traceback = _extract_stack()


### Step 6: Visualize the embeddings

In [None]:
def plot_with_labels(low_dim_embs, labels):
    assert low_dim_embs.shape[0] >= len(labels), "More labels than embeddings"
    pylab.figure(figsize=(18, 18))  # in inches
    for i, label in enumerate(labels):
        x, y = low_dim_embs[i, :]
        pylab.scatter(x, y)
        pylab.annotate(label,
                     xy=(x, y),
                     xytext=(5, 2),
                     textcoords='offset points',
                     ha='right',
                     va='bottom')

    pylab.show()
    
from sklearn.manifold import TSNE

tsne = TSNE(perplexity=30, n_components=2, init='pca', n_iter=5000)
plot_only = 500
low_dim_embs = tsne.fit_transform(final_embeddings[:plot_only, :])
labels = [reverse_dictionary[i] for i in xrange(plot_only)]
plot_with_labels(low_dim_embs, labels)