资讯动态

Flink 实战:HBase Connector 用法全解析——从环境搭建到 TableInputFormat 批量读与实时写入

发布时间:2026/10/5 16:04:06 来源:尧图企业网站定制
示例工程大数据【免费下载链接】flink-learningflink 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 namehbase.rootdir/name !-- 配置HBase存储文件的目录 -- valuefile:///usr/local/var/hbase/value /property property namehbase.zookeeper.property.clientPort/name value2181/value /property property namehbase.zookeeper.property.dataDir/name !-- 配置HBase存储内建zookeeper文件的目录 -- value/usr/local/var/zookeeper/value /property property namehbase.zookeeper.dns.interface/name valuelo0/value /property property namehbase.regionserver.dns.interface/name valuelo0/value /property property namehbase.master.dns.interface/name valuelo0/value /property /configuration这里的关键配置项含义如下hbase.rootdirHBase 存储数据的根目录单机环境下使用file://本地路径即可hbase.zookeeper.property.clientPortZooKeeper 客户端端口默认 2181Flink Connector 连接 HBase 时也需要使用该端口hbase.zookeeper.property.dataDirHBase 内建 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 进程zhishengzhisheng /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 groupIdorg.apache.flink/groupId artifactIdflink-hbase_${scala.binary.version}/artifactId version${flink.version}/version /dependency dependency groupIdorg.apache.hadoop/groupId artifactIdhadoop-common/artifactId version2.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 TableInputFormatTuple2String, String() { private Tuple2String, String reuse new Tuple2String, 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 Tuple2String, 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 FilterFunctionTuple2String, String() { Override public boolean filter(Tuple2String, 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()返回要读取的表名zhishengmapResultToTuple(Result)将 HBase 的一行Result映射为Tuple2String, 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实现OutputFormatString接口并通过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.quorumlocalhost:2181 hbase.client.retries.number1 hbase.master.info.port-1 hbase.zookeeper.property.clientPort2081 hbase.rpc.timeout30000 hbase.client.operation.timeout30000 hbase.client.scanner.timeout.period30000 # HBase table name hbase.table.namezhisheng_stream hbase.column.nameinfo_stream其中hbase.table.namezhisheng_stream和hbase.column.nameinfo_stream分别指定了写入目标表与列族。这些参数最终由 ExecutionEnvUtil.java 中的ParameterTool统一加载ParameterTool会依次从 classpath 下的application.properties、命令行参数和系统属性中读取配置因此既可以直接修改配置文件也可以在提交 Job 时用--hbase.zookeeper.quorum xxx的方式覆盖。OutputFormat 生命周期configure、open、writeRecord、closeOutputFormat接口定义了四个生命周期方法HBaseStreamWriteMain内部的HBaseOutputFormat逐一实现了它们private static class HBaseOutputFormat implements OutputFormatString { 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()完成写入。其关键代码如下DataStreamSourceString data env.addSource(new FlinkKafkaConsumer( parameterTool.get(METRICS_TOPIC), //这个 kafka topic 需要和上面的工具类的 topic 一致 new SimpleStringSchema(), props)); data.map(new MapFunctionString, 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 与列值是否正确落表。小结与反思本节围绕 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-learningflink 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点击查看免费下载相关推荐Apache Iceberg Flink Connector 实战免建 Catalog用 connectoriceberg 直接建表读写Apache Iceberg Flink Connector 实战免建 Catalog用 connectoriceberg 直接建表读写 本文围绕数据湖大数据数据存储百元内做出会语音交互的机器狗ESP-HI 组装指南xiaozhi-esp32百元内做出会语音交互的机器狗ESP HI 组装指南xiaozhi esp32 xiaozhi esp32 是一套基于 MCP 协议的开源 AI 聊天机器人人工智能大模型语音交互助手嵌入式物联网智能硬件MCP 服务TDengine Flink Connector 实战使用 Sink 与 Table Sink 将 Flink 流批数据写入 TDengineTDengine Flink Connector 实战使用 Sink 与 Table Sink 将 Flink 流批数据写入 TDengine Apache数据库时序数据库物联网大数据实时分析云原生创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价 →
↑