PySpark 기반 LDA 구현: 대규모 스트 주제 분석

대규모 사용자 리뷰 데이터에서 잠재적 주제를 추출하기 위해 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)

태그: pyspark LDA gensim jieba yarn

10월 1일 22:03에 게시됨