☰
Flink 实战:HBase Connector 用法全解析——从环境搭建到 TableInputFormat 批量读与实时写入
2026/10/5 16:03:51 网站建设 项目流程
  • 示例工程
  • 大数据

【免费下载链接】flink-learning

flink learning blog. http://www.54tianzhisheng.cn/ 含 Flink 入门、概念、原理、实战、性能调优、源码解析等内容。涉及 Flink Connector、Metrics、Library、DataStream API、Table API & SQL 等内容的学习案例,还有 Flink 落地应用的大型项目案例(PVUV、日志存储、百亿数据实时去重、监控告警)分享。欢迎大家支持我的专栏《大数据实时计算引擎 Flink 实战与性能优化》

项目地址:https://gitcode.com/gh_mirrors/fl/flink-learning
点击查看免费下载

本节内容来自《Flink 实战与性能优化》系列教程第 3.10 节,配套代码位于当前仓库的 flink-learning-connectors-hbase 模块。你将掌握 HBase 单机环境的完整搭建流程、Flink 项目接入 HBase 的依赖配置,以及通过 TableInputFormat 批量读取、通过 OutputFormat 实时写入 HBase 的两种核心编程模式,并看到仓库源码中对配置参数的实际封装方式。

HBase 与 Flink Connector 概述

HBase 是一个分布式的、面向列的开源数据库,它基于 HDFS 存储,适合海量结构化数据的随机、实时读写,因此在很多公司的数据链路中承担着在线存储与查询的职责。Flink 生态中提供了对应的 HBase Connector,使得 Flink 既可以把 HBase 当作数据源(Source)批量读取数据,也可以作为数据汇(Sink)将流式计算结果写入 HBase。本节将从零开始,讲解 HBase 环境的准备、依赖的添加,以及读取和写入两种场景的完整用法。

准备环境和依赖

在使用 Flink HBase Connector 之前,需要先在本机准备一套可用的 HBase 环境,并完成项目依赖的配置。

HBase 安装

如果你是苹果系统,可以直接使用 HomeBrew 命令安装:

brew install hbase

安装完成后,HBase 会安装到路径/usr/local/Cellar/hbase/下面,由于安装的版本不同,目录下的文件名也会不同,需要留意自己安装的具体版本号。

配置 HBase

打开libexec/conf/hbase-env.sh,修改里面的JAVA_HOME:

# The java implementation to use. Java 1.7+ required. export JAVA_HOME="/Library/Java/JavaVirtualMachines/jdk1.8.0_152.jdk/Contents/Home"

注意:这里的JAVA_HOME要根据你自己的 JDK 安装路径来配置。

接着打开libexec/conf/hbase-site.xml,配置 HBase 文件的存储目录:

<configuration> <property> <name>hbase.rootdir</name> <!-- 配置HBase存储文件的目录 --> <value>file:///usr/local/var/hbase</value> </property> <property> <name>hbase.zookeeper.property.clientPort</name> <value>2181</value> </property> <property> <name>hbase.zookeeper.property.dataDir</name> <!-- 配置HBase存储内建zookeeper文件的目录 --> <value>/usr/local/var/zookeeper</value> </property> <property> <name>hbase.zookeeper.dns.interface</name> <value>lo0</value> </property> <property> <name>hbase.regionserver.dns.interface</name> <value>lo0</value> </property> <property> <name>hbase.master.dns.interface</name> <value>lo0</value> </property> </configuration>

这里的关键配置项含义如下:

  • hbase.rootdir:HBase 存储数据的根目录,单机环境下使用file://本地路径即可;
  • hbase.zookeeper.property.clientPort:ZooKeeper 客户端端口,默认 2181,Flink Connector 连接 HBase 时也需要使用该端口;
  • hbase.zookeeper.property.dataDir:HBase 内建 ZooKeeper 元数据的存储目录;
  • 三个dns.interface配置为lo0,是 macOS 本机回环接口,用于单机模式。

运行 HBase

配置完成后,执行启动命令:

./bin/start-hbase.sh

执行后打印出来的日志类似:

starting master, logging to /usr/local/var/log/hbase/hbase-zhisheng-master-zhisheng.out

验证是否安装成功

使用jps命令查看 Java 进程:

zhisheng@zhisheng /usr/local/Cellar/hbase/1.2.9/libexec jps 91302 HMaster 62535 RemoteMavenServer 1100 91471 Jps

只要出现HMaster进程,就说明 HBase 安装运行成功。

启动 HBase Shell

执行下面命令进入 HBase 的交互式 Shell:

./bin/hbase shell

在 Shell 中可以直接执行建表、读写等操作,是验证 HBase 功能最直接的方式。

停止 HBase

需要停止服务时,执行:

./bin/stop-hbase.sh

HBase 常用命令

HBase Shell 中常用的命令包括:

  • list:列出所有已存在的表;
  • create:创建表;
  • put:写入数据;
  • get:读取单行数据;
  • scan:扫描读取数据(可读全表);
  • describe:显示表详情。

这些命令在后续准备测试数据、验证 Flink 写入结果时会频繁使用。

添加依赖

在pom.xml中添加 HBase 相关的依赖:

<dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-hbase_${scala.binary.version}</artifactId> <version>${flink.version}</version> </dependency> <dependency> <groupId>org.apache.hadoop</groupId> <artifactId>hadoop-common</artifactId> <version>2.7.4</version> </dependency>

从当前仓库的实现看,flink-learning-connectors-hbase 模块在实际落地时针对不同 HBase 版本拆分了两个子模块,并在父模块中统一引入了 Hadoop 兼容依赖:

  • flink-learning-connectors-hbase-1.4:对应 HBase 1.x 场景;
  • flink-learning-connectors-hbase-2.2:对应 HBase 2.x 场景。

两个子模块的pom.xml中实际使用的是flink-connector-hbase-2.2这个新版 Connector 依赖(参见 flink-learning-connectors-hbase-2.2/pom.xml),版本由父 POM 的${flink-connector-hbase.version}属性统一管理,并在打包阶段通过maven-shade-plugin生成可提交的 fat jar。Flink HBase Connector 中,HBase 不仅可以作为数据源,也可以作为数据汇写入数据,下面我们先来看如何从 HBase 中读取数据。

Flink 使用 TableInputFormat 读取 HBase 批量数据

这里我们使用TableInputFormat来读取 HBase 中的数据,首先准备测试数据。

准备数据

先往 HBase 中插入五条数据:

put 'zhisheng', 'first', 'info:bar', 'hello' put 'zhisheng', 'second', 'info:bar', 'zhisheng001' put 'zhisheng', 'third', 'info:bar', 'zhisheng002' put 'zhisheng', 'four', 'info:bar', 'zhisheng003' put 'zhisheng', 'five', 'info:bar', 'zhisheng004'

执行scan扫描整个zhisheng表,可以看到表中一共有五条数据,其中 rowkey 分别是first、second、third、four、five,列族为info,列名为bar。

Flink Job 代码

Flink 读取 HBase 数据的完整程序代码如下:

/** * Desc: 读取 HBase 数据 */ public class HBaseReadMain { //表名 public static final String HBASE_TABLE_NAME = "zhisheng"; // 列族 static final byte[] INFO = "info".getBytes(ConfigConstants.DEFAULT_CHARSET); //列名 static final byte[] BAR = "bar".getBytes(ConfigConstants.DEFAULT_CHARSET); public static void main(String[] args) throws Exception { ExecutionEnvironment env = ExecutionEnvironment.getExecutionEnvironment(); env.createInput(new TableInputFormat<Tuple2<String, String>>() { private Tuple2<String, String> reuse = new Tuple2<String, String>(); @Override protected Scan getScanner() { Scan scan = new Scan(); scan.addColumn(INFO, BAR); return scan; } @Override protected String getTableName() { return HBASE_TABLE_NAME; } @Override protected Tuple2<String, String> mapResultToTuple(Result result) { String key = Bytes.toString(result.getRow()); String val = Bytes.toString(result.getValue(INFO, BAR)); reuse.setField(key, 0); reuse.setField(val, 1); return reuse; } }).filter(new FilterFunction<Tuple2<String, String>>() { @Override public boolean filter(Tuple2<String, String> value) throws Exception { return value.f1.startsWith("zhisheng"); } }).print(); } }

这段代码的核心逻辑可以拆解为三层:

  1. 通过env.createInput(...)将 HBase 表作为批量输入:这里传入的是TableInputFormat的匿名子类,Flink 在批处理执行环境中会按照 InputFormat 的机制拉取 HBase 中的行数据;
  2. 实现三个关键的抽象方法:
    • getScanner():构建Scan对象,通过scan.addColumn(INFO, BAR)限定只扫描info列族下的bar列,避免全表全列扫描;
    • getTableName():返回要读取的表名zhisheng;
    • mapResultToTuple(Result):将 HBase 的一行Result映射为Tuple2<String, String>,其中 rowkey 放在第一个字段,info:bar列的值放在第二个字段。这里复用了成员变量reuse来避免频繁创建对象;
  3. 接上filter算子做流式过滤:将读取出来的全部数据过滤出 value 以zhisheng开头的记录。

运行上面的 Job 后,可以看到输出结果中已经把以zhisheng开头的四条数据(zhisheng001到zhisheng004)都打印出来了,而第一条hello因为不以zhisheng开头被过滤掉。这个例子清晰地展示了 Flink 批处理与 HBase 的接入方式:TableInputFormat负责"按表读",mapResultToTuple负责"按列取值",后续可以无缝接上任意 Flink 算子继续做计算。

Flink 使用 OutputFormat 向 HBase 写入数据

读取之外,更常见的场景是把 Flink 计算后的结果写入 HBase。当前仓库的 HBaseStreamWriteMain.java 给出了一套完整的、可直接参考的流式写入实现,其核心是自定义HBaseOutputFormat,实现OutputFormat<String>接口,并通过dataStream.writeUsingOutputFormat(new HBaseOutputFormat())将流数据落到 HBase 表中。

配置参数:从常量定义到 application.properties

先看写入所需的配置参数。仓库把 HBase 的连接参数统一收敛在 HBaseConstant.java 中,包括:

常量对应的配置 Key作用
HBASE_ZOOKEEPER_QUORUMhbase.zookeeper.quorumZooKeeper 集群地址,多个用逗号分隔
HBASE_ZOOKEEPER_PROPERTY_CLIENTPORThbase.zookeeper.property.clientPortZooKeeper 客户端端口
HBASE_RPC_TIMEOUThbase.rpc.timeoutRPC 超时时间(毫秒)
HBASE_CLIENT_OPERATION_TIMEOUThbase.client.operation.timeout客户端操作超时时间(毫秒)
HBASE_CLIENT_SCANNER_TIMEOUT_PERIODhbase.client.scanner.timeout.periodScanner 超时周期(毫秒)
HBASE_CLIENT_RETRIES_NUMBERhbase.client.retries.number客户端重试次数
HBASE_MASTER_INFO_PORThbase.master.info.portMaster 信息端口
HBASE_TABLE_NAMEhbase.table.name要写入的 HBase 表名
HBASE_COLUMN_NAMEhbase.column.name要写入的列族名

对应的默认配置文件是 application.properties,里面给出了可直接运行的示例值:

# HBase hbase.zookeeper.quorum=localhost:2181 hbase.client.retries.number=1 hbase.master.info.port=-1 hbase.zookeeper.property.clientPort=2081 hbase.rpc.timeout=30000 hbase.client.operation.timeout=30000 hbase.client.scanner.timeout.period=30000 # HBase table name hbase.table.name=zhisheng_stream hbase.column.name=info_stream

其中hbase.table.name=zhisheng_stream和hbase.column.name=info_stream分别指定了写入目标表与列族。这些参数最终由 ExecutionEnvUtil.java 中的ParameterTool统一加载:ParameterTool会依次从 classpath 下的application.properties、命令行参数和系统属性中读取配置,因此既可以直接修改配置文件,也可以在提交 Job 时用--hbase.zookeeper.quorum xxx的方式覆盖。

OutputFormat 生命周期:configure、open、writeRecord、close

OutputFormat接口定义了四个生命周期方法,HBaseStreamWriteMain内部的HBaseOutputFormat逐一实现了它们:

private static class HBaseOutputFormat implements OutputFormat<String> { private org.apache.hadoop.conf.Configuration configuration; private Connection connection = null; private String taskNumber = null; private Table table = null; private int rowNumber = 0; @Override public void configure(Configuration parameters) { configuration = HBaseConfiguration.create(); configuration.set(HBASE_ZOOKEEPER_QUORUM, ExecutionEnvUtil.PARAMETER_TOOL.get(HBASE_ZOOKEEPER_QUORUM)); configuration.set(HBASE_ZOOKEEPER_PROPERTY_CLIENTPORT, ExecutionEnvUtil.PARAMETER_TOOL.get(HBASE_ZOOKEEPER_PROPERTY_CLIENTPORT)); configuration.set(HBASE_RPC_TIMEOUT, ExecutionEnvUtil.PARAMETER_TOOL.get(HBASE_RPC_TIMEOUT)); configuration.set(HBASE_CLIENT_OPERATION_TIMEOUT, ExecutionEnvUtil.PARAMETER_TOOL.get(HBASE_CLIENT_OPERATION_TIMEOUT)); configuration.set(HBASE_CLIENT_SCANNER_TIMEOUT_PERIOD, ExecutionEnvUtil.PARAMETER_TOOL.get(HBASE_CLIENT_SCANNER_TIMEOUT_PERIOD)); } @Override public void open(int taskNumber, int numTasks) throws IOException { connection = ConnectionFactory.createConnection(configuration); TableName tableName = TableName.valueOf(ExecutionEnvUtil.PARAMETER_TOOL.get(HBASE_TABLE_NAME)); Admin admin = connection.getAdmin(); if (!admin.tableExists(tableName)) { //检查是否有该表,如果没有,创建 log.info("==============不存在表 = {}", tableName); admin.createTable(new HTableDescriptor(TableName.valueOf(ExecutionEnvUtil.PARAMETER_TOOL.get(HBASE_TABLE_NAME))) .addFamily(new HColumnDescriptor(ExecutionEnvUtil.PARAMETER_TOOL.get(HBASE_COLUMN_NAME)))); } table = connection.getTable(tableName); this.taskNumber = String.valueOf(taskNumber); } @Override public void writeRecord(String record) throws IOException { Put put = new Put(Bytes.toBytes(taskNumber + rowNumber)); put.addColumn(Bytes.toBytes(ExecutionEnvUtil.PARAMETER_TOOL.get(HBASE_COLUMN_NAME)), Bytes.toBytes("zhisheng"), Bytes.toBytes(String.valueOf(rowNumber))); rowNumber++; table.put(put); } @Override public void close() throws IOException { table.close(); connection.close(); } }

逐方法分析这套实现的设计要点:

  • configure():在算子初始化阶段调用,负责基于HBaseConfiguration.create()创建 Hadoop 配置对象,并把 ZooKeeper 地址、客户端端口、RPC 超时、操作超时、Scanner 超时等参数注入其中。这些参数与 application.properties 中的配置一一对应;
  • open():在任务启动时调用,通过ConnectionFactory.createConnection(configuration)建立真正的 HBase 连接,然后获取Admin检查目标表zhisheng_stream是否存在,不存在则用HTableDescriptor加HColumnDescriptor动态创建表并指定列族info_stream,最后通过connection.getTable(tableName)拿到可写数据的Table对象。注意Connection是重量级资源,因此放在open中创建、在close中释放,而不是每条记录都重建;
  • writeRecord():每来一条流数据调用一次。这里以taskNumber + rowNumber拼成 rowkey,将记录的编号写入列族下的zhisheng列,然后调用table.put(put)提交写入;rowNumber自增保证每个 task 内的 rowkey 不冲突;
  • close():任务结束时关闭Table和Connection,释放资源。

在 HBaseStreamWriteMain.java 的主方法中,数据源部分是一个持续产出随机数的SourceFunction(源码中保留了从 Kafka 读取metrics.topic的写法作为注释参考),最后通过dataStream.writeUsingOutputFormat(new HBaseOutputFormat())把流接到 HBase Sink 上,并执行env.execute("Flink HBase connector sink")启动作业。

另一种写法:在 MapFunction 中直连 HBase 写入

仓库中的 Main.java 提供了另一种更简单的实时写入思路:直接从 Kafka 消费数据,在MapFunction里调用writeEventToHbase()完成写入。其关键代码如下:

DataStreamSource<String> data = env.addSource(new FlinkKafkaConsumer<>( parameterTool.get(METRICS_TOPIC), //这个 kafka topic 需要和上面的工具类的 topic 一致 new SimpleStringSchema(), props)); data.map(new MapFunction<String, Object>() { @Override public Object map(String string) throws Exception { writeEventToHbase(string, parameterTool); return string; } }).print(); env.execute("flink learning connectors hbase");

writeEventToHbase()的逻辑同样是:创建HBaseConfiguration→ 设置 ZooKeeper 等参数 →ConnectionFactory.createConnection→ 检查表不存在则创建 → 以当前时间戳作为 rowkey 构造Put→table.put()写入后关闭连接。区别在于它把整个连接生命周期放在每条记录的处理函数内,写法直观,适合教学演示;而基于OutputFormat的方式复用连接、性能更优,更适合生产环境。两种写法在 flink-learning-connectors-hbase-2.2 模块中都可以直接查看完整源码。

项目运行与验证

仓库的 HBase 示例模块已经内置了完整的工程化配置:

  1. 准备环境:按上文步骤安装并启动 HBase(单机模式即可),确认jps中出现HMaster;
  2. 准备依赖:模块的pom.xml已引入flink-connector-hbase-2.2,父 POM 中声明了hadoop-common与flink-hadoop-compatibility依赖,直接mvn clean package即可完成编译打包;
  3. 修改配置:编辑 application.properties,将hbase.zookeeper.quorum、hbase.zookeeper.property.clientPort等调整为实际环境的值,hbase.table.name决定写入哪张表;
  4. 运行验证:
    • 运行读取示例,观察控制台是否按预期打印出 HBase 表内以zhisheng开头的数据;
    • 运行写入示例,然后在 HBase Shell 中执行scan 'zhisheng_stream'(或list确认表被自动创建),查看写入的 rowkey 与列值是否正确落表。

小结与反思

本节围绕 Flink HBase Connector 走通了一条完整的实战链路:从 HBase 单机环境的安装、配置、启动与常用 Shell 命令,到 Maven 依赖的添加,再到两种典型用法——用TableInputFormat配合Scan与mapResultToTuple批量读取 HBase 数据,用OutputFormat生命周期方法(configure/open/writeRecord/close)把流式数据写入 HBase。当前仓库的 flink-learning-connectors-hbase 模块还展示了面向 HBase 1.4 与 2.2 的版本拆分、统一参数常量管理以及工程化的打包配置,可以直接作为生产项目的参考骨架。需要留意的是:示例代码以教学演示为主,生产环境还应在writeRecord中引入批量缓冲(如借助BufferedMutator)以提升吞吐,并结合 Flink Checkpoint 保证写入的 Exactly-Once 语义;本文涉及的命令与配置以教程编写时的环境(macOS + HBase 单机)为背景,在 Linux 集群环境下需相应调整dns.interface与 ZooKeeper 相关配置。

  • 示例工程
  • 大数据

【免费下载链接】flink-learning

flink learning blog. http://www.54tianzhisheng.cn/ 含 Flink 入门、概念、原理、实战、性能调优、源码解析等内容。涉及 Flink Connector、Metrics、Library、DataStream API、Table API & SQL 等内容的学习案例,还有 Flink 落地应用的大型项目案例(PVUV、日志存储、百亿数据实时去重、监控告警)分享。欢迎大家支持我的专栏《大数据实时计算引擎 Flink 实战与性能优化》

项目地址:https://gitcode.com/gh_mirrors/fl/flink-learning
点击查看免费下载

相关推荐

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询