문제 및 배경
HDFS에 존재하는 데이터를 분석하여 HBase에 HFile 형식으로 대량 쓰기 작업을 수행하는 과정에서 발생하는 에러 및 그 해결 방법을 설명합니다. 이 작업은 여러 열을 한 번에 쓰는 방식으로 진행됩니다.
주요 API
RDD를 통해 HFile로 저장하는 주요 API는 다음과 같습니다:
rdd.saveAsNewAPIHadoopFile(
stagingFolder,
classOf[ImmutableBytesWritable],
classOf[KeyValue],
classOf[HFileOutputFormat2],
job.getConfiguration)
에러 및 원인
1.
HFile 경고 메세지
Does it contain files in subdirectories that correspond to column family names
이 메세지는 세 가지 이유로 인해 발생할 수 있습니다:
- 코드 구조 문제
- 데이터 소스 문제
- setMapOutputKeyClass 및 saveAsNewAPIHadoopFile 중의 Class 불일치
본 문제에서는 데이터 소스 문제로 인해 발생했습니다.
2.
순서 정의 오류
Added a key not lexically larger than previous
이 오류는 일반적인 put 작업에서는 서비스가 자동으로 정렬해주기 때문에 발생하지 않습니다. 그러나 HFile 쓰기 작업에서 이오류가 발생하면 보통은 rowkey가 정렬되지 않았다고 의심하게 됩니다. 그러나 실제 원인은 다음과 같습니다:
Spark가 HFile를 쓸 때 rowkey + 열가족 + 열명을 기반으로 정렬을 수행합니다. 따라서 데이터를 쓰기 전에 이 세 가지 요소가 모두 정렬 상태여야 합니다.
3.
Kryo 직렬화 버퍼 오버플로우
Kryo serialization failed: Buffer overflow
이 오류는 Kryo 직렬화 과정에서 데이터 크기 제한을 초과했을 때 발생합니다. 이를 해결하기 위해 다음과 같이 설정을 조정합니다:
sparkConf.set("spark.kryoserializer.buffer.max" , "256m")
sparkConf.set("spark.kryoserializer.buffer" , "64m")
데이터 전처리 및 해결 방법
데이터를 정렬하고 중복을 제거하는 과정을 거칩니다:
val key: RDD[(String, TransferTime)] = data.reduceByKey((x, y) => y)
val unitData: RDD[TransferTime] = key.map(line => line._2)
완전한 코드 예제
object BulkLoadHBase extends HBaseOperations {
def bulkLoadData(data: RDD[Any], tableName: String , columnFamily:String): Unit = {
val dataBean: RDD[TransferTime] = data.map(line => line.asInstanceOf[TransferTime])
val keyedData: RDD[(String, TransferTime)] = dataBean.map(line => (line.vintime , line))
val reducedData: RDD[(String, TransferTime)] = keyedData.reduceByKey((x, y) => y)
val filteredData: RDD[TransferTime] = reducedData.map(line => line._2)
val sortedData: RDD[TransferTime] = filteredData.sortBy(f => f.vintime)
val preparedData: RDD[(ImmutableBytesWritable, ListBuffer[KeyValue])] = sortedData.map {
line =>
val rowkey = line.vintime
val clazz = Class.forName("com.example.TransferTime")
val fields = clazz.getDeclaredFields
var columns = ListBuffer[String]()
if (fields != null && !fields.isEmpty) {
for (field <- fields) {
field.setAccessible(true)
columns.append(field.getName)
}
}
val sortedColumns = columns.sortWith(_ < _)
val ik = new ImmutableBytesWritable(Bytes.toBytes(rowkey))
var kvList = ListBuffer[KeyValue]()
for (column <- sortedColumns) {
val value = line.getClass.getDeclaredField(column).get(line).toString
val kv = new KeyValue(
Bytes.toBytes(rowkey),
Bytes.toBytes(columnFamily),
Bytes.toBytes(column),
Bytes.toBytes(value)
)
kvList.append(kv)
}
(ik, kvList)
}
val flatMappedData: RDD[(ImmutableBytesWritable, KeyValue)] = preparedData.flatMapValues {
s => s.iterator
}
val finalData: RDD[(ImmutableBytesWritable, KeyValue)] = flatMappedData.sortBy {
x => x._1
}
hfileLoad(finalData, tableName, columnFamily)
}
def hfileLoad(rdd: RDD[(ImmutableBytesWritable, KeyValue)], tableName: TableName, columnFamily: String): Unit = {
var table: Table = null
try {
val startTime = System.currentTimeMillis()
println(s"시작 시간: ${startTime}")
val stagingFolder = s"hdfs://cdh1:9000/hfile/${tableName.toString}/${System.currentTimeMillis()}"
table = connection.getTable(tableName)
if (!admin.tableExists(tableName)) {
createTable(tableName, columnFamily)
}
val job = Job.getInstance(config)
job.setJobName("HFile 로드 작업")
job.setMapOutputKeyClass(classOf[ImmutableBytesWritable])
job.setMapOutputValueClass(classOf[KeyValue])
rdd.saveAsNewAPIHadoopFile(
stagingFolder,
classOf[ImmutableBytesWritable],
classOf[KeyValue],
classOf[HFileOutputFormat2],
job.getConfiguration
)
val load = new LoadIncrementalHFiles(config)
val regionLocator = connection.getRegionLocator(tableName)
HFileOutputFormat2.configureIncrementalLoad(job, table, regionLocator)
load.doBulkLoad(new Path(stagingFolder), table.asInstanceOf[HTable])
val endTime = System.currentTimeMillis()
println(s"종료 시간: ${endTime}")
println(s"처리 시간: ${endTime - startTime}ms")
} catch {
case e: IOException => e.printStackTrace()
} finally {
if (table != null) {
try {
table.close()
} catch {
case e: IOException => e.printStackTrace()
}
}
if (connection != null) {
try {
connection.close()
} catch {
case e: IOException => e.printStackTrace()
}
}
}
}
}