참조 : 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()
2016년 5월 4일 수요일
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 절을 구현
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/
/*
다시 테이블에 저장하는 소스
*/
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/
피드 구독하기:
글 (Atom)