레이블이 Spark인 게시물을 표시합니다. 모든 게시물 표시
레이블이 Spark인 게시물을 표시합니다. 모든 게시물 표시

2016년 5월 4일 수요일

spark-hbase-connector 사용법

참조 : https://github.com/nerdammer/spark-hbase-connector

/*
 spark-shell 을 이용해야 해서 git clone 으로 프로젝트 빌드를 수행

*/

//hbase의 라이브러리를 CLASSPATH에 지정
export CLASSPATH=/usr/local/hbase/lib/*

//실행 옵션. 빌드한 spark-hbase-connector의 jar파일 지정
spark-shell --jars /usr/local/spark/spark-hbase-connector_2.10-1.0.2.jar


import it.nerdammer.spark.hbase._


/*
 hbase내의 테이블의 특정 컬럼패밀리에 특정 qualifier를 지정해서 받아오는 소스

1.타입 지정 :  [(String,...)] -> qualifier 수 만큼 있다면 차례로 들어가고 1개 더많으면 가장앞이 row key값 (테이블명)

2.select : 데이터를 받고 싶은 qualifier 이름 지정

3. inColumnFamily : 컬럼 패밀리를 지정

결과 : rdd = Array[String,String,String] 형태로 들어간다고 이해 하면 됨
*/

val rdd =sc.hbaseTable[(String,String,String)]("T_NAME").select("Q1","Q2").inColumnFamily("CF_Name")



/*
다른 컬럼패밀리의 qualifier를 받아오고 싶을때

select 에서 다른 컬럼패밀리 이름을 지정
*/

val rdd =sc.hbaseTable[(String,String,String,String)]("T_NAME").select("Q1","Q2","otherCF:Q3").inColumnFamily("CF_Name")

/*
RDD를 hbase에 저장 할때

1. toHBaseTable : 저장할 테이블 명

2. toColumns :  qualifier 명

3. inColumnFamily : 컬럼 패밀리 명
*/

rdd.toHBaseTable("T_NAME").toColumns("Q1","Q2").inColumnFamily("CF_NAME").save()



Spark-postgreSQL 연동

//Spark 실행할때 jdbc jar를 지정

spark shell --driver-class-path /usr/share/java/postgresql93-jdbc.jar

/*
DOMAIN_NAME : ip 혹은 도메인

DATABASE_NAME : 접속할 Database 이름

USER_NAME : 유저 명

PASSWORD : 비밀번호

변수를 만들어서 입력 / 하드코딩
*/
val url = "jdbc:postgresql://DOMAIN_NAME/DATABASE_NAME?user=USER_NAME&password=PASSWORD"

/*
TABLE_NAME : 데이터를 가져올 테이블 명

rows : RDD 형태로 반환 (FROM 절을 구현한것)
*/
val rows = sqlContext.load("jdbc", Map("url" -> url ,"dbtable"->"TABLE_NAME"))

/*
filter 함수 에서 WHERE 절을 구현

select 함수 에서 SELECT 절을 구현
*/
row.filter("COLUMN_NAME like 'test'").select("col1","col2")

/*
다시 테이블에 저장하는  소스
*/
import org.apache.spark.sql.Row;
import org.apache.spark.sql.types.{StructType,StructField,StringType};

/*
저장 할  테이블에 컬럼 명을 공백을 두고 생성
schema에서 StructType 타입으로 변경 됨
*/
val schemaString = "col1 col2 col3"
val schema =
  StructType(
    schemaString.split(" ").map(fieldName => StructField(fieldName, StringType, true)))

/*
데이터 (rdd) Row 형태로 변환

테스트에서는 "data1_data2_data3" 을 하나의 로우로 가진 것을 이용
*/
val rowRDD = rdd.map(_.split("_")).map(p => Row(p(0),p(1),p(2)))

/*
dfr : Row로 변환된 rdd와 schema로 Data Frame 생성 (Spark DF의 연산들을 모두 사용 가능)

insertIntoJDBC : postgreSQL에 저장하는 함수, 세번째 인자로 true를 주면 기존 테이블 삭

제 하고 다시 생성 함
*/
val dfr = sqlContext.createDataFrame(rowRDD, schema)
dfr.insertIntoJDBC(url, "TABLE_NAME", true)

참조 : https://eradiating.wordpress.com/2015/04/17/using-spark-data-sources-to-load-data-from-postgresql/