Skip to content

위키미디어 데이터 분석을 위한 전처리 #121

Description

@ppabam
  • 큰 데이터 처리 - wikimedia #110 에서 이어 집니다.

  • 위키미디어 재단에서 운영하는 위키백과를 비롯한 다양한 프로젝트는 전 세계적으로 방대한 양의 지식과 정보를 제공하며, 수많은 사용자들이 이를 활발하게 이용하고 있습니다. 이러한 위키미디어 프로젝트의 페이지뷰 데이터는 정보 소비 패턴, 사용자들의 관심사, 그리고 특정 사건이 지식 접근에 미치는 영향 등을 파악할 수 있는 매우 귀중한 자료입니다.

아래 공식 기술문서를 참고하여

1. 각 팀은 데이터 분석을 목표를 수립

  • 예) 가장 많이 조회된 페이지 식별, 페이지 조회수의 시간별 추세 분석, 모바일 및 데스크톱 트래픽 분석 ...

2. 이를 위한 전처리 작업을 진행

  • pandas jupyter 등 로컬 환경에서 처리 가능한 SUMMARY 파일(PARQUET) 생성
  • 분석 목표에 맞는 필드로 이루어진 크기가 작은 요약 파일 생성 -> 분석을 위한 활용

TIP

1. 중복 회피

  • airflow job 수행시 append 모드의 write 는 기존 데이터가 있는 경우 중복 데이터가 발생
  • 아래와 같이 overwrite 모드는 SAVE_BASE 의 기존 데이터를 모두 삭제함
  • 같은 SAVE_BASE 상태에서 데이터 중복을 회피 하고자 overwrite 모드를 사용하면서 동시에 파티션이 중복되지 않는 범위에서 데이터를 추가적으로 저장하려면 spark.conf.set("spark.sql.sources.partitionOverwriteMode", "dynamic") 을 설정하 하면 아래와 같이 저장됨
$ gcloud storage ls gs://sunsin-bucket/wiki/test/parquet
gs://sunsin-bucket/wiki/test/parquet/
gs://sunsin-bucket/wiki/test/parquet/_SUCCESS
gs://sunsin-bucket/wiki/test/parquet/date=20240101/
gs://sunsin-bucket/wiki/test/parquet/date=20240104/
from pyspark.sql import SparkSession
from pyspark.sql.functions import col, lit, input_file_name
from datetime import datetime
import sys

APP_NAME = "TestSaveParquet"
#RAW_BASE = "gs://sunsin-bucket/wiki/2024-01/dt=20240101/"
RAW_BASE = "gs://sunsin-bucket/wiki/"
SAVE_BASE = "gs://sunsin-bucket/wiki/test/parquet"
DT = sys.argv[1]

# SparkSession 생성 (마스터 URL 지정 X, spark-submit에서 설정됨)
spark = SparkSession.builder.appName(f"{APP_NAME}_{DT}").getOrCreate()

def load_pageviews(file_path, date_str):
    return spark.read.option("delimiter", " ").csv(file_path, inferSchema=True) \
        .toDF("domain", "title", "views", "size") \
        .withColumn("date", lit(date_str)) \
        .withColumn("file_path", input_file_name())

def extract_prefix_and_partition(date_str: str) -> tuple[str, str]:
    """
    "2024-01-01" → ("2024-01", "20240101")
    """
    parsed_date = datetime.strptime(date_str, "%Y-%m-%d")
    prefix = parsed_date.strftime("%Y-%m")
    partition = parsed_date.strftime("%Y%m%d")
    return prefix, partition

prefix, partition = extract_prefix_and_partition(DT)
raw_path = f"{RAW_BASE}/{prefix}/dt={partition}"
df = load_pageviews(raw_path, partition)

print("LOAD".center(33, "*"))
df.show(10)


print("SAVE START".center(33, "*"))
spark.conf.set("spark.sql.sources.partitionOverwriteMode", "dynamic")
df.write.mode("overwrite").partitionBy("date").parquet(SAVE_BASE)
print("SAVE END".center(33, "*"))

# Spark 세션 종료
spark.stop()

2. 시간 절약 전략

  • 2024년 1,2,3 월 분석을 목표로 팀원단 1개월 씩 데이터 처리
  • 위 샘플 코드를 참고하여 text 파일을 우선 스키마만 적용하여 parquet 로 저장하는 DAG 먼저 수행 ( DAG 가 돌아 가는 동안 위 공식 문서 및 일부 데이터를 확인하며 목표 수립 )
  • 목표가 수립되면 1차 전처리된 parquet 를 이용하여 분석 목표에 맞는 요약 데이터 생성

3. 공통 미션

  • 위 샘플 코드는 원천 TEXT 에서 시간 정보는 없음
  • 날짜를 파일명에서 유추 하듯 시간 정보도 함께 추출하여 dataframe 에 저장

Metadata

Metadata

Assignees

No one assigned

    Labels

    No labels
    No labels

    Projects

    No projects

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions