ML 파이프라인

ML 파이프라인 (ML Pipelines)

이 섹션에서는 ML 파이프라인(Machine Learning Pipelines) 개념을 소개해요. ML 파이프라인은 DataFrames 위에 구축된 통일된 고수준 API 세트를 제공해서, 사용자가 실용적인 머신러닝 파이프라인을 만들고 튜닝하기 쉽게 도와줘요.

출처: ML Pipelines

본문

  • 파이프라인의 주요 개념
  • DataFrame
  • 파이프라인 구성 요소
    • Transformer
    • Estimator
    • 파이프라인 구성 요소의 속성
  • Pipeline
    • 동작 방식
    • 세부 사항
  • Parameter
  • ML 영속화: 파이프라인 저장과 불러오기
    • ML 영속화의 하위 호환성
  • 코드 예제
    • 예제: Estimator, Transformer, Param
    • 예제: Pipeline
  • 모델 선택(하이퍼파라미터 튜닝)

파이프라인의 주요 개념

MLlib은 머신러닝 알고리즘의 API를 표준화해서 여러 알고리즘을 하나의 파이프라인이나 워크플로로 결합하기 쉬워지게 해요. 이 섹션은 파이프라인 API가 도입한 핵심 개념을 다뤄요. 파이프라인 개념은 대부분 scikit-learn 프로젝트에서 영감을 받았어요.

  • DataFrame: 이 ML API는 Spark SQL의 DataFrame을 ML 데이터셋으로 사용해요. 다양한 데이터 타입을 담을 수 있죠. 예를 들어 DataFrame은 텍스트, 피처 벡터, 실제 레이블, 예측값을 각각 다른 컬럼으로 가질 수 있어요.
  • Transformer: Transformer는 하나의 DataFrame을 다른 DataFrame으로 변환할 수 있는 알고리즘이에요. 예를 들어 ML 모델은 피처가 있는 DataFrame을 예측값이 있는 DataFrame으로 변환하는 Transformer예요.
  • Estimator: EstimatorDataFrame에 피팅해서 Transformer를 만들어내는 알고리즘이에요. 예를 들어 학습 알고리즘은 DataFrame에서 훈련해 모델을 만들어내는 Estimator예요.
  • Pipeline: Pipeline은 여러 TransformerEstimator를 연결해 하나의 ML 워크플로를 지정해요.
  • Parameter: 모든 TransformerEstimator는 이제 파라미터를 지정하는 공통 API를 공유해요.

DataFrame

머신러닝은 벡터, 텍스트, 이미지, 구조화된 데이터 같은 다양한 데이터 타입에 적용될 수 있어요. 이 API는 다양한 데이터 타입을 지원하기 위해 Spark SQL의 DataFrame을 채택해요.

DataFrame은 많은 기본 타입과 구조화 타입을 지원해요. 지원 타입 목록은 Spark SQL 데이터타입 참조를 보세요. Spark SQL 가이드에 나열된 타입 외에도 DataFrame은 ML Vector 타입을 사용할 수 있어요.

DataFrame은 일반 RDD로부터 암시적으로 또는 명시적으로 만들 수 있어요. 예제는 아래 코드와 Spark SQL 프로그래밍 가이드를 보세요.

DataFrame의 컬럼에는 이름이 있어요. 아래 코드 예제에서는 "text", "features", "label" 같은 이름을 사용해요.

파이프라인 구성 요소

Transformer

Transformer는 피처 변환기와 학습된 모델을 포함하는 추상화예요. 기술적으로 Transformertransform() 메서드를 구현하는데, 보통 하나 이상의 컬럼을 추가하는 방식으로 하나의 DataFrame을 다른 DataFrame으로 변환해요. 예를 들어:

  • 피처 변환기는 DataFrame을 받아 컬럼(예: text)을 읽고, 그것을 새 컬럼(예: 피처 벡터)으로 매핑한 뒤, 매핑된 컬럼이 추가된 새 DataFrame을 출력할 수 있어요.
  • 학습 모델은 DataFrame을 받아 피처 벡터를 담고 있는 컬럼을 읽고, 각 피처 벡터의 레이블을 예측한 뒤, 예측된 레이블이 컬럼으로 추가된 새 DataFrame을 출력할 수 있어요.

Estimator

Estimator는 학습 알고리즘, 또는 데이터에 피팅하거나 훈련하는 모든 알고리즘의 개념을 추상화해요. 기술적으로 Estimatorfit() 메서드를 구현하는데, DataFrame을 받아 Model(즉 Transformer)을 만들어내요. 예를 들어 LogisticRegression 같은 학습 알고리즘은 Estimator이고, fit()을 호출하면 LogisticRegressionModel을 훈련시키는데, 이 모델은 Model이므로 곧 Transformer예요.

파이프라인 구성 요소의 속성

Transformer.transform()Estimator.fit()은 모두 무상태(stateless)예요. 미래에는 대체 개념을 통해 상태 있는 알고리즘을 지원할 수도 있어요.

TransformerEstimator의 각 인스턴스는 고유한 ID를 가져요. 이 ID는 파라미터를 지정할 때 유용해요(아래에서 설명).

Pipeline

머신러닝에서는 데이터를 처리하고 학습하기 위해 일련의 알고리즘을 순서대로 실행하는 것이 흔해요. 예를 들어 간단한 텍스트 문서 처리 워크플로는 몇 가지 단계를 포함할 수 있어요.

  • 각 문서의 텍스트를 단어로 분리하기
  • 각 문서의 단어를 수치 피처 벡터로 변환하기
  • 피처 벡터와 레이블로 예측 모델 학습하기

MLlib은 이런 워크플로를 Pipeline으로 표현해요. Pipeline은 특정 순서로 실행되는 PipelineStage(TransformerEstimator)의 시퀀스로 구성돼요. 이 섹션에서는 이 간단한 워크플로를 계속 사용하는 예로 쓸 거예요.

동작 방식

Pipeline은 단계들의 시퀀스로 지정되며, 각 단계는 Transformer 또는 Estimator예요. 이 단계들은 순서대로 실행되고, 입력 DataFrame은 각 단계를 지나가며 변환돼요. Transformer 단계에서는 DataFrametransform() 메서드가 호출돼요. Estimator 단계에서는 fit() 메서드가 호출되어 Transformer(이것은 PipelineModel 또는 피팅된 Pipeline의 일부가 돼요)를 만들어내고, 그 Transformertransform() 메서드가 DataFrame에 호출돼요.

간단한 텍스트 문서 워크플로로 이를 설명할게요. 아래 그림은 Pipeline훈련 시간 사용을 나타내요.

위 그림에서 위쪽 행은 세 개의 단계가 있는 Pipeline을 나타내요. 처음 두 단계(TokenizerHashingTF)는 Transformer(파란색)이고, 세 번째(LogisticRegression)는 Estimator(빨간색)예요. 아래쪽 행은 파이프라인을 흐르는 데이터를 나타내는데, 원통 모양이 DataFrame을 의미해요. 원시 텍스트 문서와 레이블이 있는 원본 DataFramePipeline.fit() 메서드가 호출돼요. Tokenizer.transform() 메서드는 원시 텍스트 문서를 단어로 분리해 단어가 담긴 새 컬럼을 DataFrame에 추가해요. HashingTF.transform() 메서드는 단어 컬럼을 피처 벡터로 변환해 그 벡터가 담긴 새 컬럼을 DataFrame에 추가해요. 이제 LogisticRegressionEstimator이므로, Pipeline은 먼저 LogisticRegression.fit()을 호출해 LogisticRegressionModel을 만들어요. 만약 PipelineEstimator가 더 있었다면, DataFrame을 다음 단계로 넘기기 전에 LogisticRegressionModeltransform() 메서드를 DataFrame에 호출했을 거예요.

PipelineEstimator예요. 따라서 Pipelinefit() 메서드가 실행되면 PipelineModel을 만들어내는데, 이것은 Transformer예요. 이 PipelineModel테스트 시간에 사용돼요. 아래 그림이 이 사용을 보여줘요.

위 그림에서 PipelineModel은 원래 Pipeline과 같은 수의 단계를 가지지만, 원래 Pipeline의 모든 EstimatorTransformer가 됐어요. 테스트 데이터셋에 PipelineModeltransform() 메서드가 호출되면, 데이터는 피팅된 파이프라인을 순서대로 통과해요. 각 단계의 transform() 메서드는 데이터셋을 갱신하고 다음 단계로 넘겨요.

PipelinePipelineModel은 훈련 데이터와 테스트 데이터가 동일한 피처 처리 단계를 거치도록 보장해 줘요.

세부 사항

DAG Pipeline: Pipeline의 단계는 순서 있는 배열로 지정돼요. 여기서 주어진 예제는 모두 선형 Pipeline, 즉 각 단계가 이전 단계가 만든 데이터를 사용하는 Pipeline이에요. 데이터 흐름 그래프가 방향성 비순환 그래프(DAG)를 이룬다면 비선형 Pipeline을 만드는 것도 가능해요. 이 그래프는 현재 각 단계의 입력·출력 컬럼 이름(보통 파라미터로 지정)을 기반으로 암시적으로 지정돼요. Pipeline이 DAG를 이루면 단계들은 위상 순서(topological order)로 지정돼야 해요.

런타임 검사: Pipeline은 다양한 타입의 DataFrame을 다룰 수 있으므로 컴파일 타임 타입 검사를 쓸 수 없어요. PipelinePipelineModel은 대신 Pipeline을 실제로 실행하기 전에 런타임 검사를 해요. 이 타입 검사는 DataFrame스키마, 즉 DataFrame 컬럼의 데이터 타입을 설명하는 정보를 이용해 수행돼요.

고유한 Pipeline 단계: Pipeline의 단계는 고유한 인스턴스여야 해요. 예를 들어 같은 인스턴스 myHashingTFPipeline에 두 번 삽입되면 안 돼요. Pipeline 단계는 고유한 ID를 가져야 하기 때문이에요. 하지만 myHashingTF1myHashingTF2처럼 다른 인스턴스(둘 다 HashingTF 타입)는 서로 다른 ID로 만들어지므로 같은 Pipeline에 넣을 수 있어요.

Parameter

MLlib EstimatorTransformer는 파라미터를 지정하는 통일된 API를 사용해요.

Param은 자체 문서를 가진 이름 있는 파라미터예요. ParamMap은 (파라미터, 값) 쌍의 집합이에요.

알고리즘에 파라미터를 전달하는 방법은 크게 두 가지예요.

  • 인스턴스에 파라미터 설정. 예를 들어 lrLogisticRegression 인스턴스라면 lr.setMaxIter(10)을 호출해 lr.fit()이 최대 10회 반복하도록 할 수 있어요. 이 API는 spark.mllib 패키지에서 쓰던 API와 비슷해요.
  • fit()이나 transform()ParamMap 전달. ParamMap에 있는 어떤 파라미터든 setter 메서드로 이전에 지정한 파라미터를 덮어써요.

파라미터는 EstimatorTransformer의 특정 인스턴스에 속해요. 예를 들어 lr1lr2라는 두 LogisticRegression 인스턴스가 있다면, 두 maxIter 파라미터를 모두 지정한 ParamMap을 만들 수 있어요: ParamMap(lr1.maxIter -> 10, lr2.maxIter -> 20). 이는 PipelinemaxIter 파라미터가 있는 알고리즘이 두 개 있을 때 유용해요.

ML 영속화: 파이프라인 저장과 불러오기

모델이나 파이프라인을 나중에 쓰기 위해 디스크에 저장하는 것은 종종 가치가 있어요. Spark 1.6에서 파이프라인 API에 모델 가져오기/내보내기 기능이 추가됐어요. Spark 2.3부터는 spark.mlpyspark.ml의 DataFrame 기반 API가 완전히 지원돼요.

ML 영속화는 Scala, Java, Python에서 동작해요. 하지만 R은 현재 수정된 형식을 사용하므로, R에서 저장한 모델은 R에서만 다시 불러올 수 있어요. 이건 미래에 고쳐질 예정이며 SPARK-15572에서 추적돼요.

ML 영속화의 하위 호환성

일반적으로 MLlib은 ML 영속화에 대해 하위 호환성을 유지해요. 즉, 한 Spark 버전에서 ML 모델이나 파이프라인을 저장하면 미래 버전의 Spark에서 다시 불러와 사용할 수 있어야 해요. 다만 아래에서 설명하는 드문 예외가 있어요.

모델 영속화: Spark 버전 X에서 Apache Spark ML 영속화로 저장한 모델이나 파이프라인을 Spark 버전 Y에서 불러올 수 있나요?

  • 주요(major) 버전: 보장은 없지만 최선을 다해 지원해요.
  • 부(minor)·패치(patch) 버전: 예, 하위 호환돼요.
  • 형식에 대한 참고: 안정적인 영속화 형식에 대한 보장은 없지만, 모델 로딩 자체는 하위 호환되도록 설계됐어요.

모델 동작: Spark 버전 X의 모델이나 파이프라인이 Spark 버전 Y에서 동일하게 동작하나요?

  • 주요 버전: 보장은 없지만 최선을 다해 지원해요.
  • 부·패치 버전: 버그 수정을 제외하면 동일하게 동작해요.

모델 영속화와 모델 동작 모두에서, 부 버전이나 패치 버전을 넘나드는 호환성 깨짐(breaking change)은 Spark 버전 릴리스 노트에 보고돼요. 릴리스 노트에 보고되지 않은 깨짐은 고쳐야 할 버그로 취급해야 해요.

코드 예제

이 섹션은 위에서 설명한 기능을 보여주는 코드 예제를 제공해요. 자세한 내용은 API 문서(Python, Scala, Java)를 참고하세요.

예제: Estimator, Transformer, Param

이 예제는 Estimator, Transformer, Param의 개념을 다뤄요.

API에 대한 자세한 내용은 Estimator Python 문서, Transformer Python 문서, Params Python 문서를 참고하세요.

from pyspark.ml.linalg import Vectors
from pyspark.ml.classification import LogisticRegression

# Prepare training data from a list of (label, features) tuples.
training = spark.createDataFrame([
    (1.0, Vectors.dense([0.0, 1.1, 0.1])),
    (0.0, Vectors.dense([2.0, 1.0, -1.0])),
    (0.0, Vectors.dense([2.0, 1.3, 1.0])),
    (1.0, Vectors.dense([0.0, 1.2, -0.5]))], ["label", "features"])

# Create a LogisticRegression instance. This instance is an Estimator.
lr = LogisticRegression(maxIter=10, regParam=0.01)
# Print out the parameters, documentation, and any default values.
print("LogisticRegression parameters:\n" + lr.explainParams() + "\n")

# Learn a LogisticRegression model. This uses the parameters stored in lr.
model1 = lr.fit(training)

# Since model1 is a Model (i.e., a transformer produced by an Estimator),
# we can view the parameters it used during fit().
# This prints the parameter (name: value) pairs, where names are unique IDs for this
# LogisticRegression instance.
print("Model 1 was fit using parameters: ")
print(model1.extractParamMap())

# We may alternatively specify parameters using a Python dictionary as a paramMap
paramMap = {lr.maxIter: 20}
paramMap[lr.maxIter] = 30  # Specify 1 Param, overwriting the original maxIter.
# Specify multiple Params.
paramMap.update({lr.regParam: 0.1, lr.threshold: 0.55})  # type: ignore

# You can combine paramMaps, which are python dictionaries.
# Change output column name
paramMap2 = {lr.probabilityCol: "myProbability"}
paramMapCombined = paramMap.copy()
paramMapCombined.update(paramMap2)  # type: ignore

# Now learn a new model using the paramMapCombined parameters.
# paramMapCombined overrides all parameters set earlier via lr.set* methods.
model2 = lr.fit(training, paramMapCombined)
print("Model 2 was fit using parameters: ")
print(model2.extractParamMap())

# Prepare test data
test = spark.createDataFrame([
    (1.0, Vectors.dense([-1.0, 1.5, 1.3])),
    (0.0, Vectors.dense([3.0, 2.0, -0.1])),
    (1.0, Vectors.dense([0.0, 2.2, -1.5]))], ["label", "features"])

# Make predictions on test data using the Transformer.transform() method.
# LogisticRegression.transform will only use the 'features' column.
# Note that model2.transform() outputs a "myProbability" column instead of the usual
# 'probability' column since we renamed the lr.probabilityCol parameter previously.
prediction = model2.transform(test)
result = prediction.select("features", "label", "myProbability", "prediction") \
    .collect()

for row in result:
    print("features=%s, label=%s -> prob=%s, prediction=%s"
          % (row.features, row.label, row.myProbability, row.prediction))

전체 예제 코드는 Spark 저장소의 "examples/src/main/python/ml/estimator_transformer_param_example.py"에서 찾을 수 있어요.

API에 대한 자세한 내용은 Estimator Scala 문서, Transformer Scala 문서, Params Scala 문서를 참고하세요.

import org.apache.spark.ml.classification.LogisticRegression
import org.apache.spark.ml.linalg.{Vector, Vectors}
import org.apache.spark.ml.param.ParamMap
import org.apache.spark.sql.Row

// Prepare training data from a list of (label, features) tuples.
val training = spark.createDataFrame(Seq(
  (1.0, Vectors.dense(0.0, 1.1, 0.1)),
  (0.0, Vectors.dense(2.0, 1.0, -1.0)),
  (0.0, Vectors.dense(2.0, 1.3, 1.0)),
  (1.0, Vectors.dense(0.0, 1.2, -0.5))
)).toDF("label", "features")

// Create a LogisticRegression instance. This instance is an Estimator.
val lr = new LogisticRegression()
// Print out the parameters, documentation, and any default values.
println(s"LogisticRegression parameters:\n ${lr.explainParams()}\n")

// We may set parameters using setter methods.
lr.setMaxIter(10)
  .setRegParam(0.01)

// Learn a LogisticRegression model. This uses the parameters stored in lr.
val model1 = lr.fit(training)
// Since model1 is a Model (i.e., a Transformer produced by an Estimator),
// we can view the parameters it used during fit().
// This prints the parameter (name: value) pairs, where names are unique IDs for this
// LogisticRegression instance.
println(s"Model 1 was fit using parameters: ${model1.parent.extractParamMap()}")

// We may alternatively specify parameters using a ParamMap,
// which supports several methods for specifying parameters.
val paramMap = ParamMap(lr.maxIter -> 20)
  .put(lr.maxIter, 30)  // Specify 1 Param. This overwrites the original maxIter.
  .put(lr.regParam -> 0.1, lr.threshold -> 0.55)  // Specify multiple Params.

// One can also combine ParamMaps.
val paramMap2 = ParamMap(lr.probabilityCol -> "myProbability")  // Change output column name.
val paramMapCombined = paramMap ++ paramMap2

// Now learn a new model using the paramMapCombined parameters.
// paramMapCombined overrides all parameters set earlier via lr.set* methods.
val model2 = lr.fit(training, paramMapCombined)
println(s"Model 2 was fit using parameters: ${model2.parent.extractParamMap()}")

// Prepare test data.
val test = spark.createDataFrame(Seq(
  (1.0, Vectors.dense(-1.0, 1.5, 1.3)),
  (0.0, Vectors.dense(3.0, 2.0, -0.1)),
  (1.0, Vectors.dense(0.0, 2.2, -1.5))
)).toDF("label", "features")

// Make predictions on test data using the Transformer.transform() method.
// LogisticRegression.transform will only use the 'features' column.
// Note that model2.transform() outputs a 'myProbability' column instead of the usual
// 'probability' column since we renamed the lr.probabilityCol parameter previously.
model2.transform(test)
  .select("features", "label", "myProbability", "prediction")
  .collect()
  .foreach { case Row(features: Vector, label: Double, prob: Vector, prediction: Double) =>
    println(s"($features, $label) -> prob=$prob, prediction=$prediction")
  }

전체 예제 코드는 Spark 저장소의 "examples/src/main/scala/org/apache/spark/examples/ml/EstimatorTransformerParamExample.scala"에서 찾을 수 있어요.

API에 대한 자세한 내용은 Estimator Java 문서, Transformer Java 문서, Params Java 문서를 참고하세요.

import java.util.Arrays;
import java.util.List;

import org.apache.spark.ml.classification.LogisticRegression;
import org.apache.spark.ml.classification.LogisticRegressionModel;
import org.apache.spark.ml.linalg.VectorUDT;
import org.apache.spark.ml.linalg.Vectors;
import org.apache.spark.ml.param.ParamMap;
import org.apache.spark.sql.Dataset;
import org.apache.spark.sql.Row;
import org.apache.spark.sql.RowFactory;
import org.apache.spark.sql.types.DataTypes;
import org.apache.spark.sql.types.Metadata;
import org.apache.spark.sql.types.StructField;
import org.apache.spark.sql.types.StructType;

// Prepare training data.
List<Row> dataTraining = Arrays.asList(
    RowFactory.create(1.0, Vectors.dense(0.0, 1.1, 0.1)),
    RowFactory.create(0.0, Vectors.dense(2.0, 1.0, -1.0)),
    RowFactory.create(0.0, Vectors.dense(2.0, 1.3, 1.0)),
    RowFactory.create(1.0, Vectors.dense(0.0, 1.2, -0.5))
);
StructType schema = new StructType(new StructField[]{
    new StructField("label", DataTypes.DoubleType, false, Metadata.empty()),
    new StructField("features", new VectorUDT(), false, Metadata.empty())
});
Dataset<Row> training = spark.createDataFrame(dataTraining, schema);

// Create a LogisticRegression instance. This instance is an Estimator.
LogisticRegression lr = new LogisticRegression();
// Print out the parameters, documentation, and any default values.
System.out.println("LogisticRegression parameters:\n" + lr.explainParams() + "\n");

// We may set parameters using setter methods.
lr.setMaxIter(10).setRegParam(0.01);

// Learn a LogisticRegression model. This uses the parameters stored in lr.
LogisticRegressionModel model1 = lr.fit(training);
// Since model1 is a Model (i.e., a Transformer produced by an Estimator),
// we can view the parameters it used during fit().
// This prints the parameter (name: value) pairs, where names are unique IDs for this
// LogisticRegression instance.
System.out.println("Model 1 was fit using parameters: " + model1.parent().extractParamMap());

// We may alternatively specify parameters using a ParamMap.
ParamMap paramMap = new ParamMap()
  .put(lr.maxIter().w(20))  // Specify 1 Param.
  .put(lr.maxIter(), 30)  // This overwrites the original maxIter.
  .put(lr.regParam().w(0.1), lr.threshold().w(0.55));  // Specify multiple Params.

// One can also combine ParamMaps.
ParamMap paramMap2 = new ParamMap()
  .put(lr.probabilityCol().w("myProbability"));  // Change output column name
ParamMap paramMapCombined = paramMap.$plus$plus(paramMap2);

// Now learn a new model using the paramMapCombined parameters.
// paramMapCombined overrides all parameters set earlier via lr.set* methods.
LogisticRegressionModel model2 = lr.fit(training, paramMapCombined);
System.out.println("Model 2 was fit using parameters: " + model2.parent().extractParamMap());

// Prepare test documents.
List<Row> dataTest = Arrays.asList(
    RowFactory.create(1.0, Vectors.dense(-1.0, 1.5, 1.3)),
    RowFactory.create(0.0, Vectors.dense(3.0, 2.0, -0.1)),
    RowFactory.create(1.0, Vectors.dense(0.0, 2.2, -1.5))
);
Dataset<Row> test = spark.createDataFrame(dataTest, schema);

// Make predictions on test documents using the Transformer.transform() method.
// LogisticRegression.transform will only use the 'features' column.
// Note that model2.transform() outputs a 'myProbability' column instead of the usual
// 'probability' column since we renamed the lr.probabilityCol parameter previously.
Dataset<Row> results = model2.transform(test);
Dataset<Row> rows = results.select("features", "label", "myProbability", "prediction");
for (Row r: rows.collectAsList()) {
  System.out.println("(" + r.get(0) + ", " + r.get(1) + ") -> prob=" + r.get(2)
    + ", prediction=" + r.get(3));
}

전체 예제 코드는 Spark 저장소의 "examples/src/main/java/org/apache/spark/examples/ml/JavaEstimatorTransformerParamExample.java"에서 찾을 수 있어요.

예제: Pipeline

이 예제는 위 그림에서 설명한 간단한 텍스트 문서 Pipeline을 따라가요.

API에 대한 자세한 내용은 Pipeline Python 문서를 참고하세요.

from pyspark.ml import Pipeline
from pyspark.ml.classification import LogisticRegression
from pyspark.ml.feature import HashingTF, Tokenizer

# Prepare training documents from a list of (id, text, label) tuples.
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)
], ["id", "text", "label"])

# Configure an ML pipeline, which consists of three stages: tokenizer, hashingTF, and lr.
tokenizer = Tokenizer(inputCol="text", outputCol="words")
hashingTF = HashingTF(inputCol=tokenizer.getOutputCol(), outputCol="features")
lr = LogisticRegression(maxIter=10, regParam=0.001)
pipeline = Pipeline(stages=[tokenizer, hashingTF, lr])

# Fit the pipeline to training documents.
model = pipeline.fit(training)

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

# Make predictions on test documents and print columns of interest.
prediction = model.transform(test)
selected = prediction.select("id", "text", "probability", "prediction")
for row in selected.collect():
    rid, text, prob, prediction = row
    print(
        "(%d, %s) --> prob=%s, prediction=%f" % (
            rid, text, str(prob), prediction
        )
    )

전체 예제 코드는 Spark 저장소의 "examples/src/main/python/ml/pipeline_example.py"에서 찾을 수 있어요.

API에 대한 자세한 내용은 Pipeline Scala 문서를 참고하세요.

import org.apache.spark.ml.{Pipeline, PipelineModel}
import org.apache.spark.ml.classification.LogisticRegression
import org.apache.spark.ml.feature.{HashingTF, Tokenizer}
import org.apache.spark.ml.linalg.Vector
import org.apache.spark.sql.Row

// Prepare training documents from a list of (id, text, label) tuples.
val training = spark.createDataFrame(Seq(
  (0L, "a b c d e spark", 1.0),
  (1L, "b d", 0.0),
  (2L, "spark f g h", 1.0),
  (3L, "hadoop mapreduce", 0.0)
)).toDF("id", "text", "label")

// Configure an ML pipeline, which consists of three stages: tokenizer, hashingTF, and lr.
val tokenizer = new Tokenizer()
  .setInputCol("text")
  .setOutputCol("words")
val hashingTF = new HashingTF()
  .setNumFeatures(1000)
  .setInputCol(tokenizer.getOutputCol)
  .setOutputCol("features")
val lr = new LogisticRegression()
  .setMaxIter(10)
  .setRegParam(0.001)
val pipeline = new Pipeline()
  .setStages(Array(tokenizer, hashingTF, lr))

// Fit the pipeline to training documents.
val model = pipeline.fit(training)

// Now we can optionally save the fitted pipeline to disk
model.write.overwrite().save("/tmp/spark-logistic-regression-model")

// We can also save this unfit pipeline to disk
pipeline.write.overwrite().save("/tmp/unfit-lr-model")

// And load it back in during production
val sameModel = PipelineModel.load("/tmp/spark-logistic-regression-model")

// Prepare test documents, which are unlabeled (id, text) tuples.
val test = spark.createDataFrame(Seq(
  (4L, "spark i j k"),
  (5L, "l m n"),
  (6L, "spark hadoop spark"),
  (7L, "apache hadoop")
)).toDF("id", "text")

// Make predictions on test documents.
model.transform(test)
  .select("id", "text", "probability", "prediction")
  .collect()
  .foreach { case Row(id: Long, text: String, prob: Vector, prediction: Double) =>
    println(s"($id, $text) --> prob=$prob, prediction=$prediction")
  }

전체 예제 코드는 Spark 저장소의 "examples/src/main/scala/org/apache/spark/examples/ml/PipelineExample.scala"에서 찾을 수 있어요.

API에 대한 자세한 내용은 Pipeline Java 문서를 참고하세요.

import java.util.Arrays;

import org.apache.spark.ml.Pipeline;
import org.apache.spark.ml.PipelineModel;
import org.apache.spark.ml.PipelineStage;
import org.apache.spark.ml.classification.LogisticRegression;
import org.apache.spark.ml.feature.HashingTF;
import org.apache.spark.ml.feature.Tokenizer;
import org.apache.spark.sql.Dataset;
import org.apache.spark.sql.Row;

// Prepare training documents, which are labeled.
Dataset<Row> training = spark.createDataFrame(Arrays.asList(
  new JavaLabeledDocument(0L, "a b c d e spark", 1.0),
  new JavaLabeledDocument(1L, "b d", 0.0),
  new JavaLabeledDocument(2L, "spark f g h", 1.0),
  new JavaLabeledDocument(3L, "hadoop mapreduce", 0.0)
), JavaLabeledDocument.class);

// Configure an ML pipeline, which consists of three stages: tokenizer, hashingTF, and lr.
Tokenizer tokenizer = new Tokenizer()
  .setInputCol("text")
  .setOutputCol("words");
HashingTF hashingTF = new HashingTF()
  .setNumFeatures(1000)
  .setInputCol(tokenizer.getOutputCol())
  .setOutputCol("features");
LogisticRegression lr = new LogisticRegression()
  .setMaxIter(10)
  .setRegParam(0.001);
Pipeline pipeline = new Pipeline()
  .setStages(new PipelineStage[] {tokenizer, hashingTF, lr});

// Fit the pipeline to training documents.
PipelineModel model = pipeline.fit(training);

// Prepare test documents, which are unlabeled.
Dataset<Row> test = spark.createDataFrame(Arrays.asList(
  new JavaDocument(4L, "spark i j k"),
  new JavaDocument(5L, "l m n"),
  new JavaDocument(6L, "spark hadoop spark"),
  new JavaDocument(7L, "apache hadoop")
), JavaDocument.class);

// Make predictions on test documents.
Dataset<Row> predictions = model.transform(test);
for (Row r : predictions.select("id", "text", "probability", "prediction").collectAsList()) {
  System.out.println("(" + r.get(0) + ", " + r.get(1) + ") --> prob=" + r.get(2)
    + ", prediction=" + r.get(3));
}

전체 예제 코드는 Spark 저장소의 "examples/src/main/java/org/apache/spark/examples/ml/JavaPipelineExample.java"에서 찾을 수 있어요.

모델 선택(하이퍼파라미터 튜닝)

ML 파이프라인을 사용하는 큰 이점은 하이퍼파라미터 최적화예요. 자동 모델 선택에 대한 자세한 내용은 ML Tuning Guide을 참고하세요.

더 알아보기 (Learn more)