- 示例工程
- 大数据
【免费下载链接】flink-learning
flink learning blog. http://www.54tianzhisheng.cn/ 含 Flink 入门、概念、原理、实战、性能调优、源码解析等内容。涉及 Flink Connector、Metrics、Library、DataStream API、Table API & SQL 等内容的学习案例,还有 Flink 落地应用的大型项目案例(PVUV、日志存储、百亿数据实时去重、监控告警)分享。欢迎大家支持我的专栏《大数据实时计算引擎 Flink 实战与性能优化》
本节内容来自《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.shHBase 常用命令
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(); } }这段代码的核心逻辑可以拆解为三层:
- 通过
env.createInput(...)将 HBase 表作为批量输入:这里传入的是TableInputFormat的匿名子类,Flink 在批处理执行环境中会按照 InputFormat 的机制拉取 HBase 中的行数据; - 实现三个关键的抽象方法:
getScanner():构建Scan对象,通过scan.addColumn(INFO, BAR)限定只扫描info列族下的bar列,避免全表全列扫描;getTableName():返回要读取的表名zhisheng;mapResultToTuple(Result):将 HBase 的一行Result映射为Tuple2<String, String>,其中 rowkey 放在第一个字段,info:bar列的值放在第二个字段。这里复用了成员变量reuse来避免频繁创建对象;
- 接上
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_QUORUM | hbase.zookeeper.quorum | ZooKeeper 集群地址,多个用逗号分隔 |
HBASE_ZOOKEEPER_PROPERTY_CLIENTPORT | hbase.zookeeper.property.clientPort | ZooKeeper 客户端端口 |
HBASE_RPC_TIMEOUT | hbase.rpc.timeout | RPC 超时时间(毫秒) |
HBASE_CLIENT_OPERATION_TIMEOUT | hbase.client.operation.timeout | 客户端操作超时时间(毫秒) |
HBASE_CLIENT_SCANNER_TIMEOUT_PERIOD | hbase.client.scanner.timeout.period | Scanner 超时周期(毫秒) |
HBASE_CLIENT_RETRIES_NUMBER | hbase.client.retries.number | 客户端重试次数 |
HBASE_MASTER_INFO_PORT | hbase.master.info.port | Master 信息端口 |
HBASE_TABLE_NAME | hbase.table.name | 要写入的 HBase 表名 |
HBASE_COLUMN_NAME | hbase.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 示例模块已经内置了完整的工程化配置:
- 准备环境:按上文步骤安装并启动 HBase(单机模式即可),确认
jps中出现HMaster; - 准备依赖:模块的
pom.xml已引入flink-connector-hbase-2.2,父 POM 中声明了hadoop-common与flink-hadoop-compatibility依赖,直接mvn clean package即可完成编译打包; - 修改配置:编辑 application.properties,将
hbase.zookeeper.quorum、hbase.zookeeper.property.clientPort等调整为实际环境的值,hbase.table.name决定写入哪张表; - 运行验证:
- 运行读取示例,观察控制台是否按预期打印出 HBase 表内以
zhisheng开头的数据; - 运行写入示例,然后在 HBase Shell 中执行
scan 'zhisheng_stream'(或list确认表被自动创建),查看写入的 rowkey 与列值是否正确落表。
- 运行读取示例,观察控制台是否按预期打印出 HBase 表内以
小结与反思
本节围绕 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 实战与性能优化》
相关推荐
Apache Iceberg Flink Connector 实战:免建 Catalog,用 `'connector'='iceberg'` 直接建表读写
Apache Iceberg Flink Connector 实战:免建 Catalog,用 'connector'='iceberg' 直接建表读写 本文围绕
数据湖大数据数据存储百元内做出会语音交互的机器狗:ESP-HI 组装指南(xiaozhi-esp32)
百元内做出会语音交互的机器狗:ESP HI 组装指南(xiaozhi esp32) xiaozhi esp32 是一套基于 MCP 协议的开源 AI 聊天机器人
人工智能大模型语音交互助手嵌入式物联网智能硬件MCP 服务TDengine Flink Connector 实战:使用 Sink 与 Table Sink 将 Flink 流批数据写入 TDengine
TDengine Flink Connector 实战:使用 Sink 与 Table Sink 将 Flink 流批数据写入 TDengine Apache
数据库时序数据库物联网大数据实时分析云原生
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考