You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

能否使用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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.21 22:05:02