能否使用Kafka Streams写入HBase表?求Scala实现示例
Kafka Streams 写入 HBase 实现方案
当然可以用 Kafka Streams 写入 HBase。你可以通过自定义 Processor 实现灵活的数据转换(类似 DataFrame 操作)后写入 HBase,也可以用 Kafka Connect 的 HBase Sink 连接器,但自定义处理器更适配复杂转换逻辑的场景。
下面是完整的 Scala 代码示例,覆盖两种核心场景:从 Kafka Topic 读取数据转换后写入 HBase,以及从 HBase 读取数据处理后写回 HBase。
1. 依赖配置(build.sbt)
先引入必要的依赖包:
name := "kafka-streams-hbase-demo" version := "1.0" scalaVersion := "2.13.8" libraryDependencies ++= Seq( "org.apache.kafka" %% "kafka-streams" % "3.4.0", "org.apache.hbase" % "hbase-client" % "2.5.3", "org.apache.hbase" % "hbase-common" % "2.5.3", "ch.qos.logback" % "logback-classic" % "1.2.11" )
2. 自定义 HBase Sink 处理器
实现 Kafka Streams Processor,负责将处理后的数据写入 HBase:
import org.apache.hadoop.conf.Configuration import org.apache.hadoop.hbase.client.{Connection, ConnectionFactory, Put, Table, TableName} import org.apache.hadoop.hbase.util.Bytes import org.apache.kafka.streams.processor.{AbstractProcessor, ProcessorContext} class HBaseSinkProcessor(tableName: String, columnFamily: String) extends AbstractProcessor[String, String] { private var hbaseTable: Table = _ private var connection: Connection = _ override def init(context: ProcessorContext): Unit = { super.init(context) // 初始化 HBase 连接 val conf = Configuration.create() conf.set("hbase.zookeeper.quorum", "your-zookeeper-host:2181") // 替换为你的 ZK 地址 conf.set("hbase.zookeeper.property.clientPort", "2181") connection = ConnectionFactory.createConnection(conf) hbaseTable = connection.getTable(TableName.valueOf(tableName)) } override def process(key: String, value: String): Unit = { // 模拟 DataFrame 风格的字段拆分转换 val fields = value.split(",") if (fields.length >= 3) { // 构造 HBase Put 对象,用 Kafka 消息的 key 作为行键 val put = new Put(Bytes.toBytes(key)) // 写入指定列族的列 put.addColumn( Bytes.toBytes(columnFamily), Bytes.toBytes("field1"), Bytes.toBytes(fields(0)) ) put.addColumn( Bytes.toBytes(columnFamily), Bytes.toBytes("field2"), Bytes.toBytes(fields(1)) ) put.addColumn( Bytes.toBytes(columnFamily), Bytes.toBytes("field3"), Bytes.toBytes(fields(2)) ) hbaseTable.put(put) } } override def close(): Unit = { // 关闭资源 if (hbaseTable != null) hbaseTable.close() if (connection != null) connection.close() } }
3. 主程序:从 Kafka Topic 读取并写入 HBase
import org.apache.kafka.streams.{KafkaStreams, StreamsConfig, Topology} import org.apache.kafka.streams.scala.ImplicitConversions._ import org.apache.kafka.streams.scala.Serdes._ import org.apache.kafka.streams.scala.kstream.KStream import java.util.Properties object KafkaStreamsToHBase { def main(args: Array[String]): Unit = { // Kafka Streams 基础配置 val props = new Properties() props.put(StreamsConfig.APPLICATION_ID_CONFIG, "kafka-streams-hbase-demo") props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "your-kafka-broker:9092") // 替换为你的 Kafka 地址 props.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, classOf[StringSerde]) props.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, classOf[StringSerde]) val builder = new StreamsBuilder() // 从指定 Kafka Topic 读取数据 val inputStream: KStream[String, String] = builder.stream("input-topic") // 执行类似 DataFrame 的转换操作:过滤、映射、校验 val processedStream = inputStream .filter((_, value) => value.nonEmpty) // 过滤空消息 .mapValues(_.trim.toUpperCase()) // 字段值转大写 .filter((_, value) => value.split(",").length >=3) // 校验字段数量 // 绑定自定义 HBase Sink 处理器 processedStream.process(() => new HBaseSinkProcessor("hbase-target-table", "cf")) // 替换为你的 HBase 表名和列族 val topology: Topology = builder.build() val streams = new KafkaStreams(topology, props) // 启动流处理应用 streams.start() // 优雅关闭钩子 sys.addShutdownHook { streams.close() } } }
4. 从 HBase 读取数据处理后写回 HBase
先通过 HBase API 读取数据并发送到 Kafka Topic,再复用上述流处理逻辑写回 HBase:
import org.apache.hadoop.conf.Configuration import org.apache.hadoop.hbase.client.{ConnectionFactory, Scan, Table, TableName} import org.apache.hadoop.hbase.util.Bytes import org.apache.kafka.clients.producer.{KafkaProducer, ProducerRecord} import java.util.Properties object HBaseToKafkaToHBase { def main(args: Array[String]): Unit = { // 1. 从 HBase 读取数据并发送到 Kafka val hbaseConf = Configuration.create() hbaseConf.set("hbase.zookeeper.quorum", "your-zookeeper-host:2181") val connection = ConnectionFactory.createConnection(hbaseConf) val sourceTable = connection.getTable(TableName.valueOf("hbase-source-table")) val scan = new Scan() scan.addColumn(Bytes.toBytes("cf"), Bytes.toBytes("field1")) val resultScanner = sourceTable.getScanner(scan) // Kafka Producer 配置 val producerProps = new Properties() producerProps.put("bootstrap.servers", "your-kafka-broker:9092") producerProps.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer") producerProps.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer") val producer = new KafkaProducer[String, String](producerProps) // 遍历 HBase 结果并发送到 Kafka val iterator = resultScanner.iterator() while (iterator.hasNext) { val result = iterator.next() val rowKey = Bytes.toString(result.getRow) val field1 = Bytes.toString(result.getValue(Bytes.toBytes("cf"), Bytes.toBytes("field1"))) producer.send(new ProducerRecord[String, String]("hbase-input-topic", rowKey, s"$field1,processed,data")) } // 关闭资源 resultScanner.close() sourceTable.close() connection.close() producer.close() // 2. 启动流处理逻辑写回 HBase KafkaStreamsToHBase.main(args) } }
注意事项
- 替换代码中的
your-zookeeper-host、your-kafka-broker、HBase 表名和列族为实际环境值 - HBase 客户端版本需与集群版本一致,避免兼容性问题
- 生产环境建议使用连接池管理 HBase 连接,减少资源开销
- 可结合 Kafka Streams 状态存储实现更复杂的聚合、关联逻辑
内容的提问来源于stack exchange,提问作者Rohit M
相关产品推荐
相关产品推荐

