用于将文本转换为特征的标准Transformers是CountVectorizer
CountVectorizer 和 CountVectorizerModel 旨在帮助将文本文档集合转换为令牌计数向量。
或HashingTF:
使用散列技巧将一系列术语映射到它们的术语频率。目前我们使用 Austin Appleby 的 MurmurHash 3 算法(MurmurHash3_x86_32)来计算术语对象的哈希码值。由于使用简单的模将散列函数转换为列索引,因此建议使用 2 的幂作为 numFeatures 参数;否则特征将不会均匀地映射到列。
两者都有binary 选项,可用于从计数切换到二进制向量。
没有内置的 Transfomer 可以给出你想要的准确结果(它对 ML 算法没有用处)你可以购买 explode 应用 StringIndexer 和 collect_list / collect_set:
import org.apache.spark.ml.feature._
import org.apache.spark.ml.Pipeline
val df = Seq(
(1, Array("I", "like", "Spark")), (2, Array("I", "hate", "Spark"))
).toDF("id", "words")
val pipeline = new Pipeline().setStages(Array(
new SQLTransformer()
.setStatement("SELECT id, explode(words) as word FROM __THIS__"),
new StringIndexer().setInputCol("word").setOutputCol("index"),
new SQLTransformer()
.setStatement("""SELECT id, COLLECT_SET(index) AS values
FROM __THIS__ GROUP BY id""")
))
pipeline.fit(df).transform(df).show
// +---+---------------+
// | id| values|
// +---+---------------+
// | 1|[0.0, 1.0, 3.0]|
// | 2|[2.0, 0.0, 1.0]|
// +---+---------------+
使用CountVectorizer 和udf:
import org.apache.spark.ml.linalg._
spark.udf.register("indices", (v: Vector) => v.toSparse.indices)
val pipeline = new Pipeline().setStages(Array(
new CountVectorizer().setInputCol("words").setOutputCol("vector"),
new SQLTransformer()
.setStatement("SELECT *, indices(vector) FROM __THIS__")
))
pipeline.fit(df).transform(df).show
// +---+----------------+--------------------+-------------------+
// | id| words| vector|UDF:indices(vector)|
// +---+----------------+--------------------+-------------------+
// | 1|[I, like, Spark]|(4,[0,1,3],[1.0,1...| [0, 1, 3]|
// | 2|[I, hate, Spark]|(4,[0,1,2],[1.0,1...| [0, 1, 2]|
// +---+----------------+--------------------+-------------------+