/**
* 把数据写入到hbase 中
*/
object KafkaToHbase {
def main(args: Array[String]): Unit = {
//testHbase()
//创建kafka消费者的对象
val kafkaConsumer =new KafkaConsumer[String,String](PropertiesUtil.properties)
//订阅指定的topic 用于数据的消费
kafkaConsumer.subscribe(util.Arrays.asList(PropertiesUtil.getProperty("kafka.topics")))
println("等待消费数据--------------")
while (true) {
//每0.1S 从指定topic中消费数据
val records: ConsumerRecords[String,String] = kafkaConsumer.poll(100)
//这个是scala和java集合类型之间的转换
for (cr <- records) {
//得到每条数据的value
val str:String = cr.value()
println(str)
HbaseDao.put(str)
}
}
}
}