In [0]:
from pyspark.ml import Pipeline
from pyspark.ml.classification import LogisticRegression
from pyspark.ml.evaluation import BinaryClassificationEvaluator
from pyspark.ml.feature import CountVectorizer, Tokenizer
from pyspark.ml.tuning import CrossValidator, ParamGridBuilder

In [0]:
# Prepare training documents, which are labeled.
training = spark.createDataFrame([
    (0, "a b c d e spark", 1.0),
    (1, "b d", 0.0),
    (2, "spark f g h", 1.0),
    (3, "hadoop mapreduce", 0.0),
    (4, "b spark who", 1.0),
    (5, "g d a y", 0.0),
    (6, "spark fly", 1.0),
    (7, "was mapreduce", 0.0),
    (8, "e spark program", 1.0),
    (9, "a e c l", 0.0),
    (10, "spark compile", 1.0),
    (11, "hadoop software", 0.0)
], ["id", "text", "label"])

In [0]:
# Configure an ML pipeline, which consists of tree stages: tokenizer, hashingTF, and lr.
tokenizer = Tokenizer(inputCol="text", outputCol="words")
count_vectorizer = CountVectorizer(inputCol=tokenizer.getOutputCol(), outputCol="features")
lr = LogisticRegression(maxIter=10)
pipeline = Pipeline(stages=[tokenizer, count_vectorizer, lr])

In [0]:
# We use a ParamGridBuilder to construct a grid of parameters to search over.
# With 3 values for count_vectorizer.vocabSize and 2 values for lr.regParam,
# this grid will have 3 x 2 = 6 parameter settings for CrossValidator to choose from.
paramGrid = ParamGridBuilder() \
    .addGrid(count_vectorizer.vocabSize, [1024, 16384, 262144]) \
    .addGrid(lr.regParam, [0.1, 0.01]) \
    .build()

In [0]:
# A CrossValidator requires an Estimator, a set of Estimator ParamMaps, and an Evaluator.
cross_val = CrossValidator(estimator=pipeline,
                           estimatorParamMaps=paramGrid,
                           evaluator=BinaryClassificationEvaluator(),
                           numFolds=2) # use 3+ folds in practice

In [0]:
# Run cross-validation, and choose the best set of parameters.
cvModel = cross_val.fit(training)

In [0]:
print("\nModel was fit using parameters: ")
print(cvModel.extractParamMap())

In [0]:
print("\nThe best model was fit using parameters: ")
print("\nTokenizer: ")
print(cvModel.bestModel.stages[0].extractParamMap())

In [0]:
print("\nCount vectorizer: ")
print(cvModel.bestModel.stages[1].extractParamMap())

In [0]:
print("\nLogistic regression: ")
print(cvModel.bestModel.stages[2].extractParamMap())

In [0]:
# Prepare test documents, which are unlabeled.
test = spark.createDataFrame([
    (4, "spark i j k"),
    (5, "l m n"),
    (6, "mapreduce spark"),
    (7, "apache hadoop")
], ["id", "text"])

In [0]:
prediction = cvModel.transform(test)

In [0]:
# Make predictions on test documents. cvModel uses the best model found (lrModel).
prediction = cvModel.transform(test)
selected = prediction.select("id", "text", "probability", "prediction")
print('\n')
for row in selected.collect():
    print(row)
