
이번 글은 정확히는 mongoDB의 JSON 데이터를 Apache-Spark를 통해서 RDB 구조의 테이블로 ETL하는 코드에 대해서 설명할 계획이다.
참고: 이전 글(데이터 웨어하우스 구축 고군분투 이야기(1) 마지막에 나왔던 아키텍쳐의 빨간 박스 부분이다.

macOS+python으로 진행하였습니다.
import pyspark
from pyspark import SparkContext, SparkConf
from pyspark.sql import SparkSession
from pyspark.sql.functions import *
from pyspark.sql.types import *
spark를 설치하지 않아고
pip install pyspark로만 해도 spark 코드가 돌아가는 이유는 java 위에서 pyspark 라이브러리가 돌아가, 기본적인 부분들은 어느정도 커버가 가능하기 때문이다.
(spark는 시작하기전에 session을 생성해주어야 한다.)
session = SparkSession.builder\
.appName("appname")\
.config("spark.mongodb.input.uri", "mongodb+srv://<id>:<password>@---.---.mongodb.net/<db_name>.<collection_name>")\
.config("spark.mongodb.output.uri", "mongodb+srv://<id>:<password>@---.---.mongodb.net/<db_name>.<collection_name>")\
.config("spark.jars.packages", "org.mongodb.spark:mongo-spark-connector_2.12:3.0.1")\
.getOrCreate()
local에서 코드를 돌리기 위해, 굉장히 우여곡절이 많았다. 컴퓨터 환경 구제에 대해 기본 지식이 저처럼 부족하신 분들은 이 두 과정도 쉽지 않을 것...🥲
참고로 spark는 java 위에서 돌아가기 때문에, java 설치 및 버전 확인도 필요하다. 저 jars 파일이 굉장히 많이 애를 먹인다. 단순하게 말하면, 코드를 실행하는 디렉토리에 jdbc 드라이버가 존재해야한다. 그게 아니면 아래와 같은 Py4JJavaError로 시작하는 불쾌한 에러를 만난다.

이렇게 코드로 쓰면 간단해보이는데, 사내에서 db관련 권한을 받는 데에도 많은 시간이 소요되고, 환경을 구성하는데에 많은 시간이 걸린다.
mongoDB에 연결한 후, 실질적으로 데이터를 Read 하는 단계라고 볼 수 있다. ]
df = session.read.format("mongo")\
.option("uri", "mongodb+srv://<id>:<password>@---.---.mongodb.net/<db_name>.<collection_name>")\
.load()
transformation을 진행한다. 가져오는 컬렉션의 경우, column을 select만 해주어도 충분한 컬렉션이기 때문에 아래와 같이 진행했다. spark sql로 transformation이 가능하니 필요에 따라 추가 리서치를 진행하면 된다.
읽어오는 mongoDB의 데이터세트 자체가 DataFrame으로 가져와지기 때문에 DataFrame으로 바꿔줄 필요는 없다.
selected_columns = ["id", "createdAt", "updatedAt"]
df = df.select(selected_columns)
이렇게 하면 df 에 mongoDB에서 가져온 데이터를 읽은 데이터프레임이 만들어진다. (spark는 RDD라는 단위의 데이터타입을 사용하고 여기서 DataFrame은 pandas DF가 아닌 schema RDD이다.)
위 작업을 무사히 끝내면, Extract, Transformation 단계가 완료된다. 이 작업은 하나의 컬렉션에서 하나의 테이블로 데이터를 이동시키는 작업이기 때문에, 이와 같은 작업을 아주 많이 진행해야 한다... 팀에서는 중복되는 부분은 코드 모듈화를 통해 최적화를 진행할 예정이다.