HBase와 Spark
HBase와 Spark
이 문서는 HBase와 Spark 사이의 네 가지 주요 상호작용 지점을 설명해요. 기본 Spark 연동, Spark Streaming, Spark Bulk Load, SparkSQL/DataFrames를 다뤄요. 모든 예시의 코드는 실제로 실행해 볼 수 있게 정리되어 있어요.
출처: 문서
본문
Spark 자체는 이 문서의 범위 밖이에요. Spark 프로젝트와 하위 프로젝트에 대한 정보는 Spark 사이트를 참고하세요. 이 문서는 Spark와 HBase 사이의 4가지 주요 상호작용 지점에 초점을 맞춰요. 그 지점들은 다음과 같아요.
기본 Spark Spark DAG의 어느 지점에서든 HBase Connection을 가질 수 있는 능력.
Spark Streaming Spark Streaming 애플리케이션의 어느 지점에서든 HBase Connection을 가질 수 있는 능력.
Spark Bulk Load HBase에 벌크 삽입하기 위해 HBase HFiles에 직접 쓰는 능력.
SparkSQL/DataFrames HBase에 표현된 테이블을 끌어다 쓰는 SparkSQL을 작성하는 능력.
다음 섹션들에서 이 모든 상호작용 지점의 예시를 살펴볼 거예요.
기본 Spark
이 섹션은 가장 낮고 가장 간단한 수준의 Spark HBase 통합을 설명해요. 다른 모든 상호작용 지점은 여기서 설명할 개념 위에 구축돼요.
모든 Spark와 HBase 통합의 뿌리는 HBaseContext예요. HBaseContext는 HBase 구성을 받아 Spark executors로 푸시해요. 이렇게 하면 각 Spark Executor당 하나의 HBase Connection을 정적 위치에 둘 수 있어요.
참고로, Spark Executor는 Region Server와 같은 노드에 있거나 다른 노드에 있을 수 있어요. 공동 배치에 대한 의존성은 없어요. 각 Spark Executor를 다중 스레드 클라이언트 애플리케이션으로 생각하세요. 이렇게 하면 executor에서 실행되는 모든 Spark Task가 공유 Connection 객체에 접근할 수 있어요.
HBaseContext 사용 예시
이 예시는 HBaseContext를 사용해 Scala에서 RDD에 foreachPartition을 수행하는 방법을 보여줘요.
val sc = new SparkContext("local", "test")
val config = new HBaseConfiguration()
...
val hbaseContext = new HBaseContext(sc, config)
rdd.hbaseForeachPartition(hbaseContext, (it, conn) => {
val bufferedMutator = conn.getBufferedMutator(TableName.valueOf("t1"))
it.foreach((putRecord) => {
. val put = new Put(putRecord._1)
. putRecord._2.foreach((putValue) => put.addColumn(putValue._1, putValue._2, putValue._3))
. bufferedMutator.mutate(put)
})
bufferedMutator.flush()
bufferedMutator.close()
})
Java로 구현한 같은 예시는 다음과 같아요.
JavaSparkContext jsc = new JavaSparkContext(sparkConf);
try {
List<byte[]> list = new ArrayList<>();
list.add(Bytes.toBytes("1"));
...
list.add(Bytes.toBytes("5"));
JavaRDD<byte[]> rdd = jsc.parallelize(list);
Configuration conf = HBaseConfiguration.create();
JavaHBaseContext hbaseContext = new JavaHBaseContext(jsc, conf);
hbaseContext.foreachPartition(rdd,
new VoidFunction<Tuple2<Iterator<byte[]>, Connection>>() {
public void call(Tuple2<Iterator<byte[]>, Connection> t)
throws Exception {
Table table = t._2().getTable(TableName.valueOf(tableName));
BufferedMutator mutator = t._2().getBufferedMutator(TableName.valueOf(tableName));
while (t._1().hasNext()) {
byte[] b = t._1().next();
Result r = table.get(new Get(b));
if (r.getExists()) {
mutator.mutate(new Put(b));
}
}
mutator.flush();
mutator.close();
table.close();
}
});
} finally {
jsc.stop();
}
Spark와 HBase 사이의 모든 기능은 SparkSQL을 제외하고 Scala와 Java에서 모두 지원돼요. SparkSQL은 Spark가 지원하는 모든 언어를 지원해요. 이 문서의 나머지에서는 Scala 예시에 초점을 맞출 거예요.
위 예시들은 connection으로 foreachPartition을 하는 방법을 보여줘요. 다른 여러 Spark 기본 함수가 즉시 지원돼요.
bulkPut
HBase에 puts를 대규모 병렬로 보내는 것
bulkDelete
HBase에 deletes를 대규모 병렬로 보내는 것
bulkGet
새 RDD를 만들기 위해 HBase에 gets를 대규모 병렬로 보내는 것
mapPartition
HBase에 대한 완전한 접근을 위해 Connection 객체로 Spark Map 함수를 수행하는 것
hbaseRDD
RDD를 만들기 위한 분산 스캔을 단순화하는 것
이 모든 기능의 예시는 hbase-connectors 저장소의 hbase-spark integration을 참고하세요(hbase-spark 커넥터는 hbase 코어 밖에, 관련된 Apache HBase 프로젝트가 유지 관리하는 별도 저장소에 있어요).
Spark Streaming
Spark Streaming은 Spark 위에 구축된 마이크로 배치 스트림 처리 프레임워크예요. HBase와 Spark Streaming은 좋은 동반자 관계를 이루는데, HBase가 Spark Streaming과 함께 다음 혜택을 제공할 수 있기 때문이에요.
- 이동 중에 참조 데이터나 프로필 데이터를 가져올 곳
- Spark Streaming의 한 번만 처리(only once processing) 약속을 지원하는 방식으로 카운트나 집계를 저장할 곳
Spark Streaming과의 hbase-spark integration은 일반 Spark 통합 지점과 유사해요. 다음 명령들이 Spark Streaming DStream에서 바로 가능하기 때문이에요.
bulkPut
HBase에 puts를 대규모 병렬로 보내는 것
bulkDelete
HBase에 deletes를 대규모 병렬로 보내는 것
bulkGet
새 RDD를 만들기 위해 HBase에 gets를 대규모 병렬로 보내는 것
mapPartition
HBase에 대한 완전한 접근을 위해 Connection 객체로 Spark Map 함수를 수행하는 것
hbaseRDD
RDD를 만들기 위한 분산 스캔을 단순화하는 것
DStreams를 사용한 bulkPut 예시
아래는 DStreams로 bulkPut 하는 예시예요. RDD bulk put과 느낌이 매우 비슷해요.
val sc = new SparkContext("local", "test")
val config = new HBaseConfiguration()
val hbaseContext = new HBaseContext(sc, config)
val ssc = new StreamingContext(sc, Milliseconds(200))
val rdd1 = ...
val rdd2 = ...
val queue = mutable.Queue[RDD[(Array[Byte], Array[(Array[Byte],
Array[Byte], Array[Byte])])]]()
queue += rdd1
queue += rdd2
val dStream = ssc.queueStream(queue)
dStream.hbaseBulkPut(
hbaseContext,
TableName.valueOf(tableName),
(putRecord) => {
val put = new Put(putRecord._1)
putRecord._2.foreach((putValue) => put.addColumn(putValue._1, putValue._2, putValue._3))
put
})
hbaseBulkPut 함수에는 세 개의 입력이 있어요. executor의 HBase Connection에 연결되는 구성 브로드캐스트 정보를 전달하는 hbaseContext, 데이터를 넣는 테이블의 이름, 그리고 DStream의 레코드를 HBase Put 객체로 변환하는 함수입니다.
Bulk Load
Spark로 HBase에 데이터를 벌크 로드하는 두 가지 옵션이 있어요. 행에 수백만 개의 열이 있는 경우와 Spark bulk load 과정의 map 쪽 전에 열이 통합되고 파티션되지 않은 경우에 동작하는 기본 bulk load 기능이 있어요.
또한 Spark와 함께 사용하는 얇은 레코드(thin record) bulk load 옵션도 있어요. 이 두 번째 옵션은 행당 10k 미만의 열을 가진 테이블용으로 설계됐어요. 이 두 번째 옵션의 장점은 더 높은 처리량과 Spark shuffle 연산에 대한 더 적은 전체 부하예요.
두 구현 모두 MapReduce bulk load 과정과 거의 비슷하게 동작해요. partitioner가 region split에 기반해 rowkey를 파티션하고, row key가 순서대로 reducers에 보내져서 HFiles가 reduce 단계에서 직접 작성될 수 있어요.
Spark 용어로, bulk load는 Spark repartitionAndSortWithinPartitions 다음에 Spark foreachPartition으로 구현돼요.
먼저 기본 bulk load 기능을 사용하는 예시를 보겠어요.
Bulk Loading 예시
다음 예시는 Spark로 bulk loading을 보여줘요.
val sc = new SparkContext("local", "test")
val config = new HBaseConfiguration()
val hbaseContext = new HBaseContext(sc, config)
val stagingFolder = ...
val rdd = sc.parallelize(Array(
(Bytes.toBytes("1"),
(Bytes.toBytes(columnFamily1), Bytes.toBytes("a"), Bytes.toBytes("foo1"))),
(Bytes.toBytes("3"),
(Bytes.toBytes(columnFamily1), Bytes.toBytes("b"), Bytes.toBytes("foo2.b"))), ...
rdd.hbaseBulkLoad(TableName.valueOf(tableName),
t => {
val rowKey = t._1
val family:Array[Byte] = t._2(0)._1
val qualifier = t._2(0)._2
val value = t._2(0)._3
val keyFamilyQualifier= new KeyFamilyQualifier(rowKey, family, qualifier)
Seq((keyFamilyQualifier, value)).iterator
},
stagingFolder.getPath)
val load = new LoadIncrementalHFiles(config)
load.doBulkLoad(new Path(stagingFolder.getPath),
conn.getAdmin, table, conn.getRegionLocator(TableName.valueOf(tableName)))
hbaseBulkLoad 함수는 세 개의 필수 파라미터를 받아요.
- 벌크 로드할 테이블의 이름
- RDD의 레코드를 튜플 키-값 쌍으로 변환하는 함수. 튜플 키는 KeyFamilyQualifer 객체이고 값은 cell 값이에요. KeyFamilyQualifer 객체는 RowKey, Column Family, Column Qualifier를 보유해요. shuffle은 RowKey로 파티션하지만 세 값 모두로 정렬해요.
- HFile이 작성될 임시 경로
Spark bulk load 명령 다음에 HBase의 LoadIncrementalHFiles 객체를 사용해 새로 생성된 HFiles를 HBase에 로드하세요.
Spark로 벌크 로딩하는 추가 파라미터
hbaseBulkLoad에 추가 파라미터 옵션으로 다음 속성을 설정할 수 있어요.
- HFiles의 최대 파일 크기
- 컴팩션에서 HFiles를 제외하는 플래그
- 압축, bloomType, blockSize, dataBlockEncoding에 대한 Column Family 설정
추가 파라미터 사용
val sc = new SparkContext("local", "test")
val config = new HBaseConfiguration()
val hbaseContext = new HBaseContext(sc, config)
val stagingFolder = ...
val rdd = sc.parallelize(Array(
(Bytes.toBytes("1"),
(Bytes.toBytes(columnFamily1), Bytes.toBytes("a"), Bytes.toBytes("foo1"))),
(Bytes.toBytes("3"),
(Bytes.toBytes(columnFamily1), Bytes.toBytes("b"), Bytes.toBytes("foo2.b"))), ...
val familyHBaseWriterOptions = new java.util.HashMap[Array[Byte], FamilyHFileWriteOptions]
val f1Options = new FamilyHFileWriteOptions("GZ", "ROW", 128, "PREFIX")
familyHBaseWriterOptions.put(Bytes.toBytes("columnFamily1"), f1Options)
rdd.hbaseBulkLoad(TableName.valueOf(tableName),
t => {
val rowKey = t._1
val family:Array[Byte] = t._2(0)._1
val qualifier = t._2(0)._2
val value = t._2(0)._3
val keyFamilyQualifier= new KeyFamilyQualifier(rowKey, family, qualifier)
Seq((keyFamilyQualifier, value)).iterator
},
stagingFolder.getPath,
familyHBaseWriterOptions,
compactionExclude = false,
HConstants.DEFAULT_MAX_FILE_SIZE)
val load = new LoadIncrementalHFiles(config)
load.doBulkLoad(new Path(stagingFolder.getPath),
conn.getAdmin, table, conn.getRegionLocator(TableName.valueOf(tableName)))
이제 얇은 레코드(thin record) bulk load 구현을 어떻게 호출하는지 보겠어요.
얇은 레코드 bulk load 사용
val sc = new SparkContext("local", "test")
val config = new HBaseConfiguration()
val hbaseContext = new HBaseContext(sc, config)
val stagingFolder = ...
val rdd = sc.parallelize(Array(
("1",
(Bytes.toBytes(columnFamily1), Bytes.toBytes("a"), Bytes.toBytes("foo1"))),
("3",
(Bytes.toBytes(columnFamily1), Bytes.toBytes("b"), Bytes.toBytes("foo2.b"))), ...
rdd.hbaseBulkLoadThinRows(hbaseContext,
TableName.valueOf(tableName),
t => {
val rowKey = t._1
val familyQualifiersValues = new FamiliesQualifiersValues
t._2.foreach(f => {
val family:Array[Byte] = f._1
val qualifier = f._2
val value:Array[Byte] = f._3
familyQualifiersValues +=(family, qualifier, value)
})
(new ByteArrayWrapper(Bytes.toBytes(rowKey)), familyQualifiersValues)
},
stagingFolder.getPath,
new java.util.HashMap[Array[Byte], FamilyHFileWriteOptions],
compactionExclude = false,
20)
val load = new LoadIncrementalHFiles(config)
load.doBulkLoad(new Path(stagingFolder.getPath),
conn.getAdmin, table, conn.getRegionLocator(TableName.valueOf(tableName)))
얇은 행에 bulk load를 사용하는 큰 차이점은 함수가 첫 값이 row key이고 두 번째 값이 모든 column family에 대한 이 행의 모든 값을 포함하는 FamiliesQualifiersValues 객체인 튜플을 반환한다는 점이에요.
SparkSQL/DataFrames
hbase-spark integration은 Spark-1.2.0에서 도입된 DataSource API(SPARK-3247)을 활용해요. 이것은 단순한 HBase KV 저장소와 복잡한 관계형 SQL 쿼리 사이의 격차를 메워 주고, 사용자가 Spark를 사용해 HBase 위에서 복잡한 데이터 분석 작업을 수행할 수 있게 해 줘요. HBase Dataframe은 표준 Spark Dataframe이며, Hive, Orc, Parquet, JSON 등 다른 데이터 소스와 상호작용할 수 있어요.
hbase-spark integration은 파티션 프루닝(partition pruning), 열 프루닝(column pruning), 조건자 푸시다운(predicate pushdown), 데이터 지역성(data locality) 같은 핵심 기법을 적용해요.
hbase-spark integration 커넥터를 사용하려면 사용자가 HBase와 Spark 테이블 사이의 스키마 매핑을 위한 Catalog를 정의하고, 데이터를 준비해 HBase 테이블을 채운 다음, HBase DataFrame을 로드해야 해요. 그 후 사용자는 통합 쿼리를 수행하고 SQL 쿼리로 HBase 테이블의 레코드에 접근할 수 있어요. 다음은 기본 절차를 보여줘요.
Catalog 정의
def catalog = s"""{
|"table":{"namespace":"default", "name":"table1"},
|"rowkey":"key",
|"columns":{
|"col0":{"cf":"rowkey", "col":"key", "type":"string"},
|"col1":{"cf":"cf1", "col":"col1", "type":"boolean"},
|"col2":{"cf":"cf2", "col":"col2", "type":"double"},
|"col3":{"cf":"cf3", "col":"col3", "type":"float"},
|"col4":{"cf":"cf4", "col":"col4", "type":"int"},
|"col5":{"cf":"cf5", "col":"col5", "type":"bigint"},
|"col6":{"cf":"cf6", "col":"col6", "type":"smallint"},
|"col7":{"cf":"cf7", "col":"col7", "type":"string"},
|"col8":{"cf":"cf8", "col":"col8", "type":"tinyint"}
|}
|}""".stripMargin
Catalog는 HBase와 Spark 테이블 사이의 매핑을 정의해요. 이 catalog의 두 가지 중요한 부분이 있어요. 하나는 rowkey 정의이고, 다른 하나는 Spark의 테이블 열과 HBase의 column family 및 column qualifier 사이의 매핑이에요. 위는 이름이 table1이고 row key가 key이며 여러 열(col1 - col8)을 가진 HBase 테이블에 대한 스키마를 정의해요. rowkey도 특정 cf(rowkey)를 가진 열(col0)로 상세히 정의되어야 한다는 점에 주의하세요.
DataFrame 저장
case class HBaseRecord(
col0: String,
col1: Boolean,
col2: Double,
col3: Float,
col4: Int,
col5: Long,
col6: Short,
col7: String,
col8: Byte)
object HBaseRecord
{
def apply(i: Int, t: String): HBaseRecord = {
val s = s"""row${"%03d".format(i)}"""
HBaseRecord(s,
i % 2 == 0,
i.toDouble,
i.toFloat,
i,
i.toLong,
i.toShort,
s"String$i: $t",
i.toByte)
}
}
val data = (0 to 255).map { i => HBaseRecord(i, "extra")}
sc.parallelize(data).toDF.write.options(
Map(HBaseTableCatalog.tableCatalog -> catalog, HBaseTableCatalog.newTable -> "5"))
.format("org.apache.hadoop.hbase.spark ")
.save()
data는 사용자가 준비한 로컬 Scala 컬렉션이고 256개의 HBaseRecord 객체를 가져요. sc.parallelize(data) 함수는 data를 분산시켜 RDD를 형성해요. toDF는 DataFrame을 반환해요. write 함수는 DataFrame을 외부 저장 시스템(여기서는 HBase)에 쓰는 데 사용되는 DataFrameWriter를 반환해요. 지정된 스키마 catalog로 DataFrame이 주어지면 save 함수는 5개 region의 HBase 테이블을 만들고 그 안에 DataFrame을 저장해요.
DataFrame 로드
def withCatalog(cat: String): DataFrame = {
sqlContext
.read
.options(Map(HBaseTableCatalog.tableCatalog->cat))
.format("org.apache.hadoop.hbase.spark")
.load()
}
val df = withCatalog(catalog)
'withCatalog' 함수에서 sqlContext는 Spark에서 구조화된 데이터(행과 열)로 작업하기 위한 진입점인 SQLContext 변수예요. read는 데이터를 DataFrame으로 읽는 데 사용할 수 있는 DataFrameReader를 반환해요. option 함수는 기본 데이터 소스에 대한 입력 옵션을 DataFrameReader에 추가하고, format 함수는 DataFrameReader에 대한 입력 데이터 소스 형식을 지정해요. load() 함수는 입력을 DataFrame으로 로드해요. withCatalog 함수가 반환한 date frame df는 HBase 테이블에 접근하는 데 사용할 수 있어요(4.4와 4.5처럼).
Language Integrated Query
val s = df.filter(($"col0" <= "row050" && $"col0" > "row040") ||
$"col0" === "row005" ||
$"col0" <= "row005")
.select("col0", "col1", "col4")
s.show
DataFrame은 join, sort, select, filter, orderBy 등 다양한 연산을 할 수 있어요. 위 df.filter는 주어진 SQL 식을 사용해 행을 필터링해요. select는 일련의 열인 col0, col1, col4를 선택해요.
SQL 쿼리
df.registerTempTable("table1")
sqlContext.sql("select count(col1) from table1").show
registerTempTable은 테이블 이름 table1을 사용해 df DataFrame을 임시 테이블로 등록해요. 이 임시 테이블의 수명은 df를 만드는 데 사용된 SQLContext에 묶여 있어요. sqlContext.sql 함수는 사용자가 SQL 쿼리를 실행할 수 있게 해 줘요.
기타
다른 타임스탬프로 쿼리
HBaseSparkConf에서 timestamp와 관련된 네 개의 파라미터를 설정할 수 있어요. 그것들은 각각 TIMESTAMP, MIN_TIMESTAMP, MAX_TIMESTAMP, MAX_VERSIONS예요. 사용자는 MIN_TIMESTAMP와 MAX_TIMESTAMP로 다른 타임스탬프나 시간 범위를 가진 레코드를 쿼리할 수 있어요. 그 사이에 아래 예시에서 tsSpecified와 oldMs 대신 구체적인 값을 사용하세요.
아래 예시는 다른 타임스탬프로 df DataFrame을 로드하는 방법을 보여줘요. tsSpecified는 사용자가 지정해요. HBaseTableCatalog는 HBase와 Relation 관계 스키마를 정의해요. writeCatalog는 스키마 매핑을 위한 catalog를 정의해요.
val df = sqlContext.read
.options(Map(HBaseTableCatalog.tableCatalog -> writeCatalog, HBaseSparkConf.TIMESTAMP -> tsSpecified.toString))
.format("org.apache.hadoop.hbase.spark")
.load()
아래 예시는 다른 시간 범위로 df DataFrame을 로드하는 방법을 보여줘요. oldMs는 사용자가 지정해요.
val df = sqlContext.read
.options(Map(HBaseTableCatalog.tableCatalog -> writeCatalog, HBaseSparkConf.MIN_TIMESTAMP -> "0",
HBaseSparkConf.MAX_TIMESTAMP -> oldMs.toString))
.format("org.apache.hadoop.hbase.spark")
.load()
df DataFrame을 로드한 후 사용자는 데이터를 쿼리할 수 있어요.
df.registerTempTable("table")
sqlContext.sql("select count(col1) from table").show
네이티브 Avro 지원
hbase-spark integration 커넥터는 Avro, JSON 등 다양한 데이터 형식을 지원해요. 아래 사용 사례는 spark가 Avro를 어떻게 지원하는지 보여줘요. 사용자는 Avro 레코드를 HBase에 직접 영속화할 수 있어요. 내부적으로 Avro 스키마는 자동으로 네이티브 Spark Catalyst 데이터 타입으로 변환돼요. HBase 테이블의 키-값 부분이 모두 Avro 형식으로 정의될 수 있다는 점에 유의하세요.
-
스키마 매핑을 위한 catalog 정의:
def catalog = s"""{ |"table":{"namespace":"default", "name":"Avrotable"}, |"rowkey":"key", |"columns":{ |"col0":{"cf":"rowkey", "col":"key", "type":"string"}, |"col1":{"cf":"cf1", "col":"col1", "type":"binary"} |} |}""".stripMargincatalog는 이름이Avrotable인 HBase 테이블에 대한 스키마예요. row key는 key이고 하나의 열 col1이 있어요. rowkey도 특정 cf(rowkey)를 가진 열(col0)로 상세히 정의되어야 해요. -
데이터 준비:
object AvroHBaseRecord { val schemaString = s"""{"namespace": "example.avro", | "type": "record", "name": "User", | "fields": [ | {"name": "name", "type": "string"}, | {"name": "favorite_number", "type": ["int", "null"]}, | {"name": "favorite_color", "type": ["string", "null"]}, | {"name": "favorite_array", "type": {"type": "array", "items": "string"}}, | {"name": "favorite_map", "type": {"type": "map", "values": "int"}} | ] }""".stripMargin val avroSchema: Schema = { val p = new Schema.Parser p.parse(schemaString) } def apply(i: Int): AvroHBaseRecord = { val user = new GenericData.Record(avroSchema); user.put("name", s"name${"%03d".format(i)}") user.put("favorite_number", i) user.put("favorite_color", s"color${"%03d".format(i)}") val favoriteArray = new GenericData.Array[String](2, avroSchema.getField("favorite_array").schema()) favoriteArray.add(s"number${i}") favoriteArray.add(s"number${i+1}") user.put("favorite_array", favoriteArray) import collection.JavaConverters._ val favoriteMap = Map[String, Int](("key1" -> i), ("key2" -> (i+1))).asJava user.put("favorite_map", favoriteMap) val avroByte = AvroSedes.serialize(user, avroSchema) AvroHBaseRecord(s"name${"%03d".format(i)}", avroByte) } } val data = (0 to 255).map { i => AvroHBaseRecord(i) }schemaString이 먼저 정의되고, 그다음 파싱되어avroSchema를 얻어요.avroSchema는AvroHBaseRecord를 생성하는 데 사용돼요. 사용자가 준비하는data는 256개의AvroHBaseRecord객체를 가진 로컬 Scala 컬렉션이에요. -
DataFrame 저장:
sc.parallelize(data).toDF.write.options( Map(HBaseTableCatalog.tableCatalog -> catalog, HBaseTableCatalog.newTable -> "5")) .format("org.apache.spark.sql.execution.datasources.hbase") .save()지정된 스키마
catalog로 data frame이 주어지면, 위는 5개 region의 HBase 테이블을 만들고 data frame을 그 안에 저장해요. -
DataFrame 로드:
def avroCatalog = s"""{ |"table":{"namespace":"default", "name":"avrotable"}, |"rowkey":"key", |"columns":{ |"col0":{"cf":"rowkey", "col":"key", "type":"string"}, |"col1":{"cf":"cf1", "col":"col1", "avro":"avroSchema"} |} |}""".stripMargin def withCatalog(cat: String): DataFrame = { sqlContext .read .options(Map("avroSchema" -> AvroHBaseRecord.schemaString, HBaseTableCatalog.tableCatalog -> avroCatalog)) .format("org.apache.spark.sql.execution.datasources.hbase") .load() } val df = withCatalog(catalog)withCatalog함수에서read는 데이터를 DataFrame으로 읽는 데 사용할 수 있는 DataFrameReader를 반환해요.option함수는 기본 데이터 소스에 대한 입력 옵션을 DataFrameReader에 추가해요. 두 가지 옵션이 있어요. 하나는avroSchema를AvroHBaseRecord.schemaString으로 설정하는 것이고, 하나는HBaseTableCatalog.tableCatalog를avroCatalog로 설정하는 것이에요.load()함수는 입력을 DataFrame으로 로드해요.withCatalog함수가 반환한 date framedf는 HBase 테이블에 접근하는 데 사용할 수 있어요. -
SQL 쿼리:
df.registerTempTable("avrotable") val c = sqlContext.sql("select count(1) from avrotable").df DataFrame을 로드한 후 사용자는 데이터를 쿼리할 수 있어요. registerTempTable은 테이블 이름 avrotable을 사용해 df DataFrame을 임시 테이블로 등록해요.
sqlContext.sql함수는 사용자가 SQL 쿼리를 실행할 수 있게 해 줘요.