기본 통계

기본 통계 (Basic Statistics)

spark.ml에서 제공하는 기본 통계 기능을 정리한 문서예요. 두 데이터 계열 간의 상관관계(Correlation) 계산, 가설 검정(Hypothesis testing), 그리고 벡터 컬럼 요약 통계(Summarizer)를 Python·Scala·Java 예제와 함께 확인해 볼게요.

출처: 문서

본문

상관관계 (Correlation)

두 데이터 계열 간의 상관관계를 계산하는 것은 통계에서 흔한 작업이에요. spark.ml에서는 여러 계열 간의 쌍별(pairwise) 상관관계를 계산할 수 있는 유연성을 제공해요. 현재 지원되는 상관 방법은 Pearson 상관과 Spearman 상관이에요.

Correlation은 지정된 방법을 사용해 입력 Vector Dataset의 상관 행렬(correlation matrix)을 계산해요. 출력은 벡터 컬럼의 상관 행렬을 포함하는 DataFrame이에요.

from pyspark.ml.linalg import Vectors
from pyspark.ml.stat import Correlation

data = [(Vectors.sparse(4, [(0, 1.0), (3, -2.0)]),),
        (Vectors.dense([4.0, 5.0, 0.0, 3.0]),),
        (Vectors.dense([6.0, 7.0, 0.0, 8.0]),),
        (Vectors.sparse(4, [(0, 9.0), (3, 1.0)]),)]
df = spark.createDataFrame(data, ["features"])

r1 = Correlation.corr(df, "features").head()

print("Pearson correlation matrix:\n" + str(r1[0]))

r2 = Correlation.corr(df, "features", "spearman").head()

print("Spearman correlation matrix:\n" + str(r2[0]))

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

Correlation은 지정된 방법을 사용해 입력 Vector Dataset의 상관 행렬을 계산해요. 출력은 벡터 컬럼의 상관 행렬을 포함하는 DataFrame이에요.

import org.apache.spark.ml.linalg.{Matrix, Vectors}
import org.apache.spark.ml.stat.Correlation
import org.apache.spark.sql.Row

val data = Seq(
  Vectors.sparse(4, Seq((0, 1.0), (3, -2.0))),
  Vectors.dense(4.0, 5.0, 0.0, 3.0),
  Vectors.dense(6.0, 7.0, 0.0, 8.0),
  Vectors.sparse(4, Seq((0, 9.0), (3, 1.0)))
)

val df = data.map(Tuple1.apply).toDF("features")
val Row(coeff1: Matrix) = Correlation.corr(df, "features").head()
println(s"Pearson correlation matrix:\n $coeff1")

val Row(coeff2: Matrix) = Correlation.corr(df, "features", "spearman").head()
println(s"Spearman correlation matrix:\n $coeff2")

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

Correlation은 지정된 방법을 사용해 입력 Vector Dataset의 상관 행렬을 계산해요. 출력은 벡터 컬럼의 상관 행렬을 포함하는 DataFrame이에요.

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

import org.apache.spark.ml.linalg.Vectors;
import org.apache.spark.ml.linalg.VectorUDT;
import org.apache.spark.ml.stat.Correlation;
import org.apache.spark.sql.Dataset;
import org.apache.spark.sql.Row;
import org.apache.spark.sql.RowFactory;
import org.apache.spark.sql.types.*;

List<Row> data = Arrays.asList(
  RowFactory.create(Vectors.sparse(4, new int[]{0, 3}, new double[]{1.0, -2.0})),
  RowFactory.create(Vectors.dense(4.0, 5.0, 0.0, 3.0)),
  RowFactory.create(Vectors.dense(6.0, 7.0, 0.0, 8.0)),
  RowFactory.create(Vectors.sparse(4, new int[]{0, 3}, new double[]{9.0, 1.0}))
);

StructType schema = new StructType(new StructField[]{
  new StructField("features", new VectorUDT(), false, Metadata.empty()),
});

Dataset<Row> df = spark.createDataFrame(data, schema);
Row r1 = Correlation.corr(df, "features").head();
System.out.println("Pearson correlation matrix:\n" + r1.get(0).toString());

Row r2 = Correlation.corr(df, "features", "spearman").head();
System.out.println("Spearman correlation matrix:\n" + r2.get(0).toString());

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

가설 검정 (Hypothesis testing)

가설 검정(Hypothesis testing)은 어떤 결과가 통계적으로 유의미한지, 즉 그 결과가 우연히 발생한 것인지 아닌지를 판단하는 데 강력한 도구예요. spark.ml은 현재 독립성 검정을 위한 Pearson 카이제곱($\chi^2$) 검정을 지원해요.

ChiSquareTest

ChiSquareTest는 레이블에 대해 각 피처마다 Pearson 독립성 검정을 수행해요. 각 피처에 대해 (피처, 레이블) 쌍이 분할표(contingency matrix)로 변환되고, 이에 대해 카이제곱 통계량이 계산돼요. 모든 레이블과 피처 값은 범주형(categorical)이어야 해요.

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

from pyspark.ml.linalg import Vectors
from pyspark.ml.stat import ChiSquareTest

data = [(0.0, Vectors.dense(0.5, 10.0)),
        (0.0, Vectors.dense(1.5, 20.0)),
        (1.0, Vectors.dense(1.5, 30.0)),
        (0.0, Vectors.dense(3.5, 30.0)),
        (0.0, Vectors.dense(3.5, 40.0)),
        (1.0, Vectors.dense(3.5, 40.0))]
df = spark.createDataFrame(data, ["label", "features"])

r = ChiSquareTest.test(df, "features", "label").head()

print("pValues: " + str(r.pValues))
print("degreesOfFreedom: " + str(r.degreesOfFreedom))
print("statistics: " + str(r.statistics))

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

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

import org.apache.spark.ml.linalg.{Vector, Vectors}
import org.apache.spark.ml.stat.ChiSquareTest

val data = Seq(
  (0.0, Vectors.dense(0.5, 10.0)),
  (0.0, Vectors.dense(1.5, 20.0)),
  (1.0, Vectors.dense(1.5, 30.0)),
  (0.0, Vectors.dense(3.5, 30.0)),
  (0.0, Vectors.dense(3.5, 40.0)),
  (1.0, Vectors.dense(3.5, 40.0))
)

val df = data.toDF("label", "features")
val chi = ChiSquareTest.test(df, "features", "label").head()
println(s"pValues = ${chi.getAs[Vector](0)}")
println(s"degreesOfFreedom ${chi.getSeq[Int](1).mkString("[", ",", "]")}")
println(s"statistics ${chi.getAs[Vector](2)}")

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

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

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

import org.apache.spark.ml.linalg.Vectors;
import org.apache.spark.ml.linalg.VectorUDT;
import org.apache.spark.ml.stat.ChiSquareTest;
import org.apache.spark.sql.Dataset;
import org.apache.spark.sql.Row;
import org.apache.spark.sql.RowFactory;
import org.apache.spark.sql.types.*;

List<Row> data = Arrays.asList(
  RowFactory.create(0.0, Vectors.dense(0.5, 10.0)),
  RowFactory.create(0.0, Vectors.dense(1.5, 20.0)),
  RowFactory.create(1.0, Vectors.dense(1.5, 30.0)),
  RowFactory.create(0.0, Vectors.dense(3.5, 30.0)),
  RowFactory.create(0.0, Vectors.dense(3.5, 40.0)),
  RowFactory.create(1.0, Vectors.dense(3.5, 40.0))
);

StructType schema = new StructType(new StructField[]{
  new StructField("label", DataTypes.DoubleType, false, Metadata.empty()),
  new StructField("features", new VectorUDT(), false, Metadata.empty()),
});

Dataset<Row> df = spark.createDataFrame(data, schema);
Row r = ChiSquareTest.test(df, "features", "label").head();
System.out.println("pValues: " + r.get(0).toString());
System.out.println("degreesOfFreedom: " + r.getList(1).toString());
System.out.println("statistics: " + r.get(2).toString());

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

Summarizer

Dataframe의 벡터 컬럼 요약 통계를 Summarizer를 통해 제공해요. 사용 가능한 지표(metric)는 컬럼별 max, min, mean, sum, variance, std, nonzeros(0이 아닌 값의 수)와 전체 count(count)예요.

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

from pyspark.ml.stat import Summarizer
from pyspark.sql import Row
from pyspark.ml.linalg import Vectors

df = sc.parallelize([Row(weight=1.0, features=Vectors.dense(1.0, 1.0, 1.0)),
                     Row(weight=0.0, features=Vectors.dense(1.0, 2.0, 3.0))]).toDF()

# "mean"과 "count" 두 지표를 위한 summarizer 생성
summarizer = Summarizer.metrics("mean", "count")

# 가중치를 사용해 여러 지표에 대한 통계 계산
df.select(summarizer.summary(df.features, df.weight)).show(truncate=False)

# 가중치 없이 여러 지표에 대한 통계 계산
df.select(summarizer.summary(df.features)).show(truncate=False)

# 가중치를 사용해 단일 지표 "mean" 계산
df.select(Summarizer.mean(df.features, df.weight)).show(truncate=False)

# 가중치 없이 단일 지표 "mean" 계산
df.select(Summarizer.mean(df.features)).show(truncate=False)

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

다음 예제는 가중치 컬럼을 사용하거나 사용하지 않고 입력 dataframe의 벡터 컬럼에 대한 mean과 variance를 계산하기 위해 Summarizer를 사용하는 방법을 보여줘요.

import org.apache.spark.ml.linalg.{Vector, Vectors}
import org.apache.spark.ml.stat.Summarizer

val data = Seq(
  (Vectors.dense(2.0, 3.0, 5.0), 1.0),
  (Vectors.dense(4.0, 6.0, 7.0), 2.0)
)

val df = data.toDF("features", "weight")

val (meanVal, varianceVal) = df.select(metrics("mean", "variance")
  .summary($"features", $"weight").as("summary"))
  .select("summary.mean", "summary.variance")
  .as[(Vector, Vector)].first()

println(s"with weight: mean = ${meanVal}, variance = ${varianceVal}")

val (meanVal2, varianceVal2) = df.select(mean($"features"), variance($"features"))
  .as[(Vector, Vector)].first()

println(s"without weight: mean = ${meanVal2}, sum = ${varianceVal2}")

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

다음 예제는 가중치 컬럼을 사용하거나 사용하지 않고 입력 dataframe의 벡터 컬럼에 대한 mean과 variance를 계산하기 위해 Summarizer를 사용하는 방법을 보여줘요.

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

import org.apache.spark.ml.linalg.Vector;
import org.apache.spark.ml.linalg.Vectors;
import org.apache.spark.ml.linalg.VectorUDT;
import org.apache.spark.ml.stat.Summarizer;
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;

List<Row> data = Arrays.asList(
  RowFactory.create(Vectors.dense(2.0, 3.0, 5.0), 1.0),
  RowFactory.create(Vectors.dense(4.0, 6.0, 7.0), 2.0)
);

StructType schema = new StructType(new StructField[]{
  new StructField("features", new VectorUDT(), false, Metadata.empty()),
  new StructField("weight", DataTypes.DoubleType, false, Metadata.empty())
});

Dataset<Row> df = spark.createDataFrame(data, schema);

Row result1 = df.select(Summarizer.metrics("mean", "variance")
  .summary(new Column("features"), new Column("weight")).as("summary"))
  .select("summary.mean", "summary.variance").first();
System.out.println("with weight: mean = " + result1.<Vector>getAs(0).toString() +
  ", variance = " + result1.<Vector>getAs(1).toString());

Row result2 = df.select(
  Summarizer.mean(new Column("features")),
  Summarizer.variance(new Column("features"))
).first();
System.out.println("without weight: mean = " + result2.<Vector>getAs(0).toString() +
  ", variance = " + result2.<Vector>getAs(1).toString());

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

더 알아보기 (Learn more)