spark on kubernetes - sparkoperator
external minio on kubernetes
minio 실행 후 pyspark로 결과 출력
참고: https://www.youtube.com/watch?v=ZzFdYm_DqEM&t=307s
helm으로 minio 실행
helm 이용해 minio 실행
helm install --set accessKey=${REDACTED},secretKey=${MASKED} --generate-name minio/minio
kubefwd로 했을 대 외부에서 접근을 못해서, nodepport를 사용함
$ minio-svc-nodeport.yaml
apiVersion: v1
kind: Service
metadata:
name: minio-svc-nodeport
spec:
ports:
- name: minio
port: 9000
targetPort: 9000
selector:
app: minio
type: NodePort
$ kubectl apply -f minio-svc-nodeport.yaml
$ kubectl get svc
NAME TYPE CLUSTER-IP EXTERNAL-IP PORT(S) AGE
minio ClusterIP 10.233.53.17 <none> 9000/TCP 4d5h
minio-svc-nodeport NodePort 10.233.37.143 <none> 9000:31867/TCP 110m
minio client는 쿠버네티스 외부에 있기 때문에, nodeport로 노출시킨 port로 통신해야함
설치 : https://docs.min.io/docs/minio-client-complete-guide.html
address-console: 9090, minio server: 9000 (헷갈리기 쉬움)
$ /.mc/config.json
{
"version": "10",
"aliases": {
"minio1": {
"url": "http://vminio01.encore.sec:9000",
"accessKey": "H11N1J3WBE5LLPX1FGPC",
"secretKey": "PyN51DsBKy3gHp8jCu4PicRsd3HUVPmhHMLsBG9c",
"api": "S3v4",
"path": "auto"
},
"minio9": {
"url": "http://172.17.172.86:31867/",
"accessKey": "myaccesskey",
"secretKey": "mysecretkey",
"api": "S3v4",
"path": "auto"
}
}
}
#bucket 만들기
$ /.mc mb minio9/test
chrome 통해 orders.json 업로드
- http://172.17.172.86:31867/
- myaccessKey, mysecretKey
- orders.json
{"id":1,"amount":1000}
{"id":2,"amount":2000}
업로드 되었는지 확인
$ ./mc ls minio9/test
[2021-09-14 13:47:27 KST] 70B orders.json
$ ./mc ls minio1/test
[2021-09-23 09:45:31 KST] 70B orders.json
pyenv, pyenv-virtualenv 세팅
sudo yum install gcc zlib-devel bzip2 bzip2-devel readline-devel sqlite \
sqlite-devel openssl-devel xz xz-devel libffi-devel
# libffi-devel가 pyenv install 보다 먼저 실행되어야함, 만약 설치가 이미되어있으면 pyenv uninstall 3.8.9, pyenv install -v 3.8.9
sudo -i
git clone https://github.com/yyuu/pyenv.git ~/.pyenv
git clone https://github.com/yyuu/pyenv-virtualenv.git ~/.pyenv/plugins/pyenv-virtualenv
환경 변수 설정
vi ~/.bashrc
(생략)
# Source global definitions
if [ -f /etc/bashrc ]; then
. /etc/bashrc
fi
export SPARK_HOME=/root/project/spark-3.1.2-bin-hadoop3.2 #spark 경로
export PATH=$PATH:$SPARK_HOME/bin
export JAVA_HOME=/home/java/jdk1.8.0_301
export PATH=$PATH:$HOME/bin:$JAVAdd_HOME/bin
export PYENV_ROOT=/root/.pyenv
export PATH=$PYENV_ROOT/bin:$PATH
eval "$(pyenv init --path)"
eval "$(pyenv virtualenv-init -)"
# 반영
$ source ~/.bashrc
pyenv로 python 환경 설정하기
pyenv versions
pyenv install -v 3.8.9
pyenv versions
pyenv global 3.8.9
출처: https://realpython.com/intro-to-pyenv/
출처: https://rmohan.com/?p=7792
spark-submit main.py
spark 설치
wget https://dlcdn.apache.org/spark/spark-3.1.2/spark-3.1.2-bin-hadoop3.2.tgz
mkdir sparkjob
cd ./sparkjob
# 파이썬 환경 3.8.9로 설정
pyenv virtualenv 3.8.9 job-3.8.9
pyenv local job-3.8.9
# pyspark 다운로드
$ cd job-3.8.9/
(job-3.8.9) root@kube01: ~/project/sparkjob/job-3.8.9:]# ll
(job-3.8.9) root@kube01: ~/project/sparkjob/job-3.8.9:] pip --trusted-host pypi.org --trusted-host files.pythonhosted.org install pyspark
python -m pip install --upgrade --trusted-host pypi.org --trusted-host files.pythonhosted.org pip
$mkdir sparkjob
$vi main.py
from pyspark import SparkContext
from pyspark.sql import SparkSession
spark = SparkSession.builder.getOrCreate()
def load_config(spark_context: SparkContext):
spark_context._jsc.hadoopConfiguration().set('fs.s3a.access.key', 'myaccesskey')
spark_context._jsc.hadoopConfiguration().set('fs.s3a.secret.key', 'mysecretkey')
spark_context._jsc.hadoopConfiguration().set('fs.s3a.path.style.access', 'true')
spark_context._jsc.hadoopConfiguration().set('fs.s3a.impl', 'org.apache.hadoop.fs.s3a.S3AFileSystem')
spark_context._jsc.hadoopConfiguration().set('fs.s3a.endpoint', 'http://172.17.172.86:31867') //외부에 공개된 port 설정
spark_context._jsc.hadoopConfiguration().set('fs.s3a.connection.ssl.enabled', 'false')
load_config(spark.sparkContext)
dataframe = spark.read.json('s3a://test/*') //orders.json 파일 선택
average = dataframe.agg({'amount': 'avg'}) //orderjs.json 파일의 amount의 평균값
average.show()
jar 가 없다고 오류가 발생
wget https://repo1.maven.org/maven2/org/apache/hadoop/hadoop-aws/3.2.0/hadoop-aws-3.2.0.jar
cp hadoop-aws-3.2.0.jar ./spark-3.1.2-bin-hadoop3.2/jars/
wget https://repo1.maven.org/maven2/com/amazonaws/aws-java-sdk-bundle/1.11.375/aws-java-sdk-bundle-1.11.375.jar
cp aws-java-sdk-bundle-1.11.375.jar ./spark-3.1.2-bin-hadoop3.2/jars/
main.py 실행
#환경 변수에 spark-submit 잡혀있음
$ spark-submit main.py
+-----------+
|avg(amount)|
+-----------+
| 2800.0|
+-----------+
21/09/14 15:47:05 INFO SparkUI: Stopped Spark web UI at http://kube01:4040
21/09/14 15:47:05 INFO MapOutputTrackerMasterEndpoint: MapOutputTrackerMasterEndpoint stopped!
minio가 잘 떠있는지 확인하기 위해, pod 안에서 통신함
alicek106/ubuntu:curl 대신 busybox를 사용해도됨
$kubectl run -i --tty --rm debug --image=vharbor.encore.sec/library/alicek106/ubuntu:curl --restart=Never bash
$ kubectl get pod -o wide
# 여기서 host는 pod의 ip, cluster가 아님
$ /bin/bash
bucket=test
file=orders.json
host=10.233.98.120:9000
s3_key='myaccesskey'
s3_secret='${REDACTED}'
resource="/${bucket}/${file}"
content_type="application/octet-stream"
date=`date -R`
_signature="GET\n\n${content_type}\n${date}\n${resource}"
signature=`echo -en ${_signature} | openssl sha1 -hmac ${s3_secret} -binary | base64`
curl -v -X GET \
-H "Host: $host" \
-H "Date: ${date}" \
-H "Content-Type: ${content_type}" \
-H "Authorization: AWS ${s3_key}:${signature}" \
http://$host${resource}
(결과)
+-----------+
|avg(amount)|
+-----------+
| 2800.0|
+-----------+