빈발 패턴 마이닝 - RDD 기반 API
빈발 패턴 마이닝 - RDD 기반 API (Frequent Pattern Mining)
대규모 데이터셋 분석의 첫 단계로 흔히 수행하는 빈발 아이템·아이템셋·부분수열 마이닝을 다루는 문서예요. spark.mllib이 제공하는 FP-growth, 연관 규칙(Association Rules), 그리고 PrefixSpan 시퀀스 패턴 마이닝을 Python·Scala·Java 예제와 함께 알아볼게요.
출처: 문서
본문
빈발 아이템, 아이템셋, 부분수열(subsequence) 또는 기타 하부 구조를 마이닝하는 것은 대규모 데이터셋을 분석하는 첫 단계 중 하나로, 수년간 데이터 마이닝의 활발한 연구 주제였어요. 자세한 내용은 Wikipedia의 association rule learning 문서를 참고하세요. spark.mllib은 빈발 아이템셋을 마이닝하는 인기 알고리즘인 FP-growth의 병렬 구현을 제공해요.
FP-growth
FP-growth 알고리즘은 Han et al., Mining frequent patterns without candidate generation 논문에 설명되어 있는데, "FP"는 빈발 패턴(frequent pattern)을 뜻해요. 트랜잭션 데이터셋이 주어졌을 때, FP-growth의 첫 단계는 아이템 빈도를 계산하고 빈발 아이템을 식별하는 것이에요. 같은 목적을 위해 설계된 Apriori류 알고리즘과 달리, FP-growth의 두 번째 단계는 후보 집합(candidate set)을 명시적으로 생성하지 않고 FP-tree라는 접미사 트리 구조로 트랜잭션을 인코딩해요. 후보 집합 생성은 보통 비용이 많이 들거든요. 두 번째 단계 이후에는 FP-tree에서 빈발 아이템셋을 추출할 수 있어요. spark.mllib에서는 PFP라고 하는 FP-growth의 병렬 버전을 구현했는데, 이는 Li et al., PFP: Parallel FP-growth for query recommendation에 설명되어 있어요. PFP는 트랜잭션의 접미사에 기반해 FP-tree를 성장시키는 작업을 분산하므로 단일 머신 구현보다 확장성이 뛰어나요. 자세한 내용은 논문들을 참고하세요.
spark.mllib의 FP-growth 구현은 다음 (하이퍼)파라미터를 사용해요:
minSupport: 아이템셋이 빈발(frequent)로 식별되기 위한 최소 지지도(support)예요. 예를 들어 어떤 아이템이 5개 트랜잭션 중 3개에 나타나면, 그 지지도는 3/5=0.6이에요.numPartitions: 작업을 분산하는 데 사용되는 파티션 수예요.
예제 (Examples)
FPGrowth가 FP-growth 알고리즘을 구현해요. 각 트랜잭션이 제네릭 타입 아이템의 List인 트랜잭션들의 RDD를 받아요. 트랜잭션으로 FPGrowth.train을 호출하면 빈발 아이템셋과 그 빈도를 저장하는 FPGrowthModel을 반환해요.
API에 대한 자세한 내용은 FPGrowth Python 문서를 참고하세요.
from pyspark.mllib.fpm import FPGrowth
data = sc.textFile("data/mllib/sample_fpgrowth.txt")
transactions = data.map(lambda line: line.strip().split(' '))
model = FPGrowth.train(transactions, minSupport=0.2, numPartitions=10)
result = model.freqItemsets().collect()
for fi in result:
print(fi)
전체 예제 코드는 Spark 저장소의 examples/src/main/python/mllib/fpgrowth_example.py 에서 찾을 수 있어요.
FPGrowth가 FP-growth 알고리즘을 구현해요. 각 트랜잭션이 제네릭 타입 아이템의 Array인 트랜잭션들의 RDD를 받아요. 트랜잭션으로 FPGrowth.run을 호출하면 빈발 아이템셋과 그 빈도를 저장하는 FPGrowthModel을 반환해요. 다음 예제는 transactions에서 빈발 아이템셋과 연관 규칙(자세한 내용은 Association Rules 참고)을 마이닝하는 방법을 보여줘요.
API에 대한 자세한 내용은 FPGrowth Scala 문서를 참고하세요.
import org.apache.spark.mllib.fpm.FPGrowth
import org.apache.spark.rdd.RDD
val data = sc.textFile("data/mllib/sample_fpgrowth.txt")
val transactions: RDD[Array[String]] = data.map(s => s.trim.split(' '))
val fpg = new FPGrowth()
.setMinSupport(0.2)
.setNumPartitions(10)
val model = fpg.run(transactions)
model.freqItemsets.collect().foreach { itemset =>
println(s"${itemset.items.mkString("[", ",", "]")},${itemset.freq}")
}
val minConfidence = 0.8
model.generateAssociationRules(minConfidence).collect().foreach { rule =>
println(s"${rule.antecedent.mkString("[", ",", "]")}=> " +
s"${rule.consequent .mkString("[", ",", "]")},${rule.confidence}")
}
전체 예제 코드는 Spark 저장소의 examples/src/main/scala/org/apache/spark/examples/mllib/SimpleFPGrowth.scala 에서 찾을 수 있어요.
FPGrowth가 FP-growth 알고리즘을 구현해요. 각 트랜잭션이 제네릭 타입 아이템의 Iterable인 트랜잭션들의 JavaRDD를 받아요. 트랜잭션으로 FPGrowth.run을 호출하면 빈발 아이템셋과 그 빈도를 저장하는 FPGrowthModel을 반환해요. 다음 예제는 transactions에서 빈발 아이템셋과 연관 규칙(자세한 내용은 Association Rules 참고)을 마이닝하는 방법을 보여줘요.
API에 대한 자세한 내용은 FPGrowth Java 문서를 참고하세요.
import java.util.Arrays;
import java.util.List;
import org.apache.spark.api.java.JavaRDD;
import org.apache.spark.api.java.JavaSparkContext;
import org.apache.spark.mllib.fpm.AssociationRules;
import org.apache.spark.mllib.fpm.FPGrowth;
import org.apache.spark.mllib.fpm.FPGrowthModel;
JavaRDD<String> data = sc.textFile("data/mllib/sample_fpgrowth.txt");
JavaRDD<List<String>> transactions = data.map(line -> Arrays.asList(line.split(" ")));
FPGrowth fpg = new FPGrowth()
.setMinSupport(0.2)
.setNumPartitions(10);
FPGrowthModel<String> model = fpg.run(transactions);
for (FPGrowth.FreqItemset<String> itemset: model.freqItemsets().toJavaRDD().collect()) {
System.out.println("[" + itemset.javaItems() + "], " + itemset.freq());
}
double minConfidence = 0.8;
for (AssociationRules.Rule<String> rule
: model.generateAssociationRules(minConfidence).toJavaRDD().collect()) {
System.out.println(
rule.javaAntecedent() + " => " + rule.javaConsequent() + ", " + rule.confidence());
}
전체 예제 코드는 Spark 저장소의 examples/src/main/java/org/apache/spark/examples/mllib/JavaSimpleFPGrowth.java 에서 찾을 수 있어요.
연관 규칙 (Association Rules)
AssociationRules는 결과(consequent)로 단일 아이템을 가지는 규칙을 구성하기 위한 병렬 규칙 생성 알고리즘을 구현해요.
API에 대한 자세한 내용은 AssociationRules Scala 문서를 참고하세요.
import org.apache.spark.mllib.fpm.AssociationRules
import org.apache.spark.mllib.fpm.FPGrowth.FreqItemset
val freqItemsets = sc.parallelize(Seq(
new FreqItemset(Array("a"), 15L),
new FreqItemset(Array("b"), 35L),
new FreqItemset(Array("a", "b"), 12L)
))
val ar = new AssociationRules()
.setMinConfidence(0.8)
val results = ar.run(freqItemsets)
results.collect().foreach { rule =>
println(s"[${rule.antecedent.mkString(",")}=>${rule.consequent.mkString(",")} ]" +
s" ${rule.confidence}")
}
전체 예제 코드는 Spark 저장소의 examples/src/main/scala/org/apache/spark/examples/mllib/AssociationRulesExample.scala 에서 찾을 수 있어요.
AssociationRules는 결과(consequent)로 단일 아이템을 가지는 규칙을 구성하기 위한 병렬 규칙 생성 알고리즘을 구현해요.
API에 대한 자세한 내용은 AssociationRules Java 문서를 참고하세요.
import java.util.Arrays;
import org.apache.spark.api.java.JavaRDD;
import org.apache.spark.api.java.JavaSparkContext;
import org.apache.spark.mllib.fpm.AssociationRules;
import org.apache.spark.mllib.fpm.FPGrowth;
import org.apache.spark.mllib.fpm.FPGrowth.FreqItemset;
JavaRDD<FPGrowth.FreqItemset<String>> freqItemsets = sc.parallelize(Arrays.asList(
new FreqItemset<>(new String[] {"a"}, 15L),
new FreqItemset<>(new String[] {"b"}, 35L),
new FreqItemset<>(new String[] {"a", "b"}, 12L)
));
AssociationRules arules = new AssociationRules()
.setMinConfidence(0.8);
JavaRDD<AssociationRules.Rule<String>> results = arules.run(freqItemsets);
for (AssociationRules.Rule<String> rule : results.collect()) {
System.out.println(
rule.javaAntecedent() + " => " + rule.javaConsequent() + ", " + rule.confidence());
}
전체 예제 코드는 Spark 저장소의 examples/src/main/java/org/apache/spark/examples/mllib/JavaAssociationRulesExample.java 에서 찾을 수 있어요.
PrefixSpan
PrefixSpan은 Pei et al., Mining Sequential Patterns by Pattern-Growth: The PrefixSpan Approach에 설명된 시퀀스 패턴 마이닝(sequential pattern mining) 알고리즘이에요. 시퀀스 패턴 마이닝 문제의 정형화는 참조된 논문을 참고하세요.
spark.mllib의 PrefixSpan 구현은 다음 파라미터를 사용해요:
minSupport: 빈발 시퀀스 패턴으로 간주되기 위한 최소 지지도예요.maxPatternLength: 빈발 시퀀스 패턴의 최대 길이예요. 이 길이를 초과하는 빈발 패턴은 결과에 포함되지 않아요.maxLocalProjDBSize: 사영(프로젝션)된 데이터베이스의 로컬 반복 처리가 시작되기 전에, 접두사-사영 데이터베이스(prefix-projected database)에 허용되는 최대 아이템 수예요. 이 파라미터는 여러분의 executor 크기에 맞춰 튜닝해야 해요.
예제 (Examples)
다음 예제는 PrefixSpan이 시퀀스들에서 동작하는 것을 보여줘요 (Pei et al.과 같은 표기법 사용):
<(12)3>
<1(32)(12)>
<(12)5>
<6>
PrefixSpan이 PrefixSpan 알고리즘을 구현해요. PrefixSpan.run을 호출하면 빈발 시퀀스와 그 빈도를 저장하는 PrefixSpanModel을 반환해요.
API에 대한 자세한 내용은 PrefixSpan Scala 문서와 PrefixSpanModel Scala 문서를 참고하세요.
import org.apache.spark.mllib.fpm.PrefixSpan
val sequences = sc.parallelize(Seq(
Array(Array(1, 2), Array(3)),
Array(Array(1), Array(3, 2), Array(1, 2)),
Array(Array(1, 2), Array(5)),
Array(Array(6))
), 2).cache()
val prefixSpan = new PrefixSpan()
.setMinSupport(0.5)
.setMaxPatternLength(5)
val model = prefixSpan.run(sequences)
model.freqSequences.collect().foreach { freqSequence =>
println(
s"${freqSequence.sequence.map(_.mkString("[", ", ", "]")).mkString("[", ", ", "]")}," +
s" ${freqSequence.freq}")
}
전체 예제 코드는 Spark 저장소의 examples/src/main/scala/org/apache/spark/examples/mllib/PrefixSpanExample.scala 에서 찾을 수 있어요.
PrefixSpan이 PrefixSpan 알고리즘을 구현해요. PrefixSpan.run을 호출하면 빈발 시퀀스와 그 빈도를 저장하는 PrefixSpanModel을 반환해요.
API에 대한 자세한 내용은 PrefixSpan Java 문서와 PrefixSpanModel Java 문서를 참고하세요.
import java.util.Arrays;
import java.util.List;
import org.apache.spark.mllib.fpm.PrefixSpan;
import org.apache.spark.mllib.fpm.PrefixSpanModel;
JavaRDD<List<List<Integer>>> sequences = sc.parallelize(Arrays.asList(
Arrays.asList(Arrays.asList(1, 2), Arrays.asList(3)),
Arrays.asList(Arrays.asList(1), Arrays.asList(3, 2), Arrays.asList(1, 2)),
Arrays.asList(Arrays.asList(1, 2), Arrays.asList(5)),
Arrays.asList(Arrays.asList(6))
), 2);
PrefixSpan prefixSpan = new PrefixSpan()
.setMinSupport(0.5)
.setMaxPatternLength(5);
PrefixSpanModel<Integer> model = prefixSpan.run(sequences);
for (PrefixSpan.FreqSequence<Integer> freqSeq: model.freqSequences().toJavaRDD().collect()) {
System.out.println(freqSeq.javaSequence() + ", " + freqSeq.freq());
}
전체 예제 코드는 Spark 저장소의 examples/src/main/java/org/apache/spark/examples/mllib/JavaPrefixSpanExample.java 에서 찾을 수 있어요.
더 알아보기 (Learn more)
- 아파치 스파크 빈발 패턴 마이닝 (원문)
- MLlib 가이드 (원문) — MLlib 전체 가이드
- ML 빈발 패턴 마이닝 — DataFrame 기반 빈발 패턴 마이닝