대규모 사용자 리뷰 데이터에서 잠재적 주제를 추출하기 위해 PySpark와 LDA(Latent Dirichlet Allocation)를 결합한 파이프라인을 구축한다. 이 글에서는 분산 환경에서의 텍스트 전처리, 모델 훈련, 결과 해석 전 과정을 다룬다.
SparkSession 초기화
YARN 클스터에 연결하고 Hive 메타스토어를 활용하는 세션을 생성한다. 원격 Python 환경 경로를 명시하여 노드 간 실행 환경을 통일한다.
from pyspark.sql import SparkSession
sess = SparkSession.builder \
.appName("LdaTopicMining") \
.master("yarn") \
.config("spark.pyspark.python", "/opt/apps/anaconda3/envs/lda_env/bin/python") \
.config("spark.sql.warehouse.dir", "/hive/warehouse") \
.config("hive.metastore.uris", "thrift://master01:9083") \
.config("spark.sql.parquet.writeLegacyFormat", "true") \
.enableHiveSupport() \
.getOrCreate()
텍스트 정제 파이프라인
jieba로 형태소 분석하고, NLTK 불용어 사전과 사용자 정의 필터로 노이즈를 제거한다. 정규식으로 특수문자를 걸러낸 뒤 재분석하는 방식으로 파이프라인을 구성한다.
import jieba
import re
from nltk.corpus import stopwords
class ReviewPreprocessor:
def __init__(self):
self.stop_words = set(stopwords.words('chinese'))
self.punct_pat = re.compile(
r'[!?。。"#$%&'()*+,-/:;<=>@[\]^_`{|}~'
r'⦅⦆「」、、〃》「」『』【】〔〕〖〗〘〙〚〛〜〝〞〟〰〾〿–—‘’‛""„‟…‧﹏'
r'.\t\n\s]+'
)
def segment(self, raw_text):
return list(jieba.cut(raw_text))
def filter_tokens(self, token_list):
filtered = []
for tok in token_list:
if tok not in self.stop_words and len(tok) > 1:
cleaned = self.punct_pat.sub('', tok)
if cleaned:
filtered.append(cleaned)
return filtered
def process(self, document):
tokens = self.segment(document)
tokens = self.filter_tokens(tokens)
return ' '.join(tokens)
분산 환경 LDA 모델링
Spark DataFrame에서 추출한 텍스트를 RDD로 변환한 뒤, gensim의 Dictionary와 LdaModel로 주제를 추출한다. 단일 문서에 대해 토픽 수를 1로 설정하고 50회 반복하여 수렴시킨다.
from gensim import corpora, models
def extract_topics(text_blob, top_n=8):
analyzer = ReviewPreprocessor()
words = analyzer.process(text_blob).split()
if len(words) < 3:
return []
vocab = corpora.Dictionary([words])
bow = [vocab.doc2bow(words)]
lda = models.LdaModel(
bow,
num_topics=1,
id2word=vocab,
passes=50,
random_state=42
)
return lda.show_topics(num_words=top_n, formatted=False)
Hive 테이블 데이터 처리
JSON 형태로 저장된 리뷰 목록을 파싱하여 의미 있는 내용만 합친 , 주제 분석을 적용한다. 기본 평가 문구는 필터링하여 제외한다.
import json
def parse_reviews(json_blob):
try:
records = json.loads(str(json_blob), strict=False)
except (json.JSONDecodeError, TypeError):
return None
combined = []
for batch in records:
for entry in batch:
body = entry.get("content", "")
if body != "用户未点评,系统默认好评。":
combined.append(body)
return "".join(combined) if combined else None
def analyze_dataset(table_name="cjw_data.qvna", sample_size=100):
dataset = sess.table(table_name)
rows = dataset.take(sample_size)
results = []
for row in rows:
raw = row[-1]
merged = parse_reviews(raw)
if merged:
topics = extract_topics(merged)
results.append((merged[:120], topics))
for topic_id, terms in topics:
print(f"[Topic {topic_id}] {terms}")
return results
짧은 스트 대응 전략
LDA는 충분한 문맥을 요구하므로 짧은 리에 대해서는 다음 접근을 고려한다.
- 문서 집계: 동일 사용자 또는 동일 상품의 리뷰를 하나의 문서로 병합
- 파라미터 튜닝:
num_topics를 낮추고passes를 높여 안정성 확보 - 대안 모델: NMF나 LSA로 대체하거나 TF-IDF 가중치를 도입
- 확 임계값: 토픽 할당 확률이 낮으면 "미분류"로 처리
통합 실행 스크립트
import jieba
from gensim import corpora, models
import nltk
import json
import re
from pyspark.sql import SparkSession
nltk.download('stopwords', quiet=True)
if __name__ == "__main__":
sess = SparkSession.builder \
.appName("LdaTopicMining") \
.master("yarn") \
.config("spark.pyspark.python", "/opt/apps/anaconda3/envs/lda_env/bin/python") \
.config("spark.sql.warehouse.dir", "/hive/warehouse") \
.config("hive.metastore.uris", "thrift://master01:9083") \
.config("spark.sql.parquet.writeLegacyFormat", "true") \
.enableHiveSupport() \
.getOrCreate()
analyze_dataset("cjw_data.qvna", 100)