话单分析之KafkaToHbase

/**

* 把数据写入到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)

}

}

}

}

©著作权归作者所有,转载或内容合作请联系作者
【社区内容提示】社区部分内容疑似由AI辅助生成,浏览时请结合常识与多方信息审慎甄别。
平台声明:文章内容(如有图片或视频亦包括在内)由作者上传并发布,文章内容仅代表作者本人观点,简书系信息发布平台,仅提供信息存储服务。

友情链接更多精彩内容