Spark 대량로드(HFile) 에러 및 해결 방법

문제 및 배경

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()
        }
      }
    }
  }
}

태그: Spark HBase Bulk Load HFile

7월 27일 13:25에 게시됨