Flink1.12.1通过Table API / Flink SQL读取HBase2.4.0
昨天群里有人问 Flink 1.12 读取Hbase的问题,于是看到这篇文章分享给大家。本文作者Ashiamd。
1. 环境
废话不多说,这里用到的环境如下(不确定是否都必要,但是至少我是这个环境)
- zookeeper 3.6.2
- Hbase 2.4.0
- Flink 1.12.1
2. HBase表
# 创建表
create 'u_m_01' , 'u_m_r'
# 插入数据
put 'u_m_01', 'a,A', 'u_m_r:r' , '1'
put 'u_m_01', 'a,B', 'u_m_r:r' , '3'
put 'u_m_01', 'b,B', 'u_m_r:r' , '3'
put 'u_m_01', 'b,C', 'u_m_r:r' , '4'
put 'u_m_01', 'c,A', 'u_m_r:r' , '2'
put 'u_m_01', 'c,C', 'u_m_r:r' , '5'
put 'u_m_01', 'c,D', 'u_m_r:r' , '1'
put 'u_m_01', 'd,B', 'u_m_r:r' , '5'
put 'u_m_01', 'd,D', 'u_m_r:r' , '2'
put 'u_m_01', 'e,A', 'u_m_r:r' , '3'
put 'u_m_01', 'e,B', 'u_m_r:r' , '2'
put 'u_m_01', 'f,A', 'u_m_r:r' , '1'
put 'u_m_01', 'f,B', 'u_m_r:r' , '2'
put 'u_m_01', 'f,D', 'u_m_r:r' , '3'
put 'u_m_01', 'g,C', 'u_m_r:r' , '1'
put 'u_m_01', 'g,D', 'u_m_r:r' , '4'
put 'u_m_01', 'h,A', 'u_m_r:r' , '1'
put 'u_m_01', 'h,B', 'u_m_r:r' , '2'
put 'u_m_01', 'h,C', 'u_m_r:r' , '4'
put 'u_m_01', 'h,D', 'u_m_r:r' , '5'
3. pom依赖
- jdk1.8
- Flink1.12.1 使用的pom依赖如下(有些是多余的)
<?xml version="1.0" encoding="UTF-8"?>
="http://maven.apache.org/POM/4.0.0"
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
>4.0.0</modelVersion>
>org.example</groupId>
>flink-hive-hbase</artifactId>
>1.0-SNAPSHOT</version>
>
.compiler.source>8</maven.compiler.source>
.compiler.target>8</maven.compiler.target>
.version>1.12.1</flink.version>
.binary.version>2.12</scala.binary.version>
.version>3.1.2</hive.version>
.version>8.0.19</mysql.version>
.version>2.4.0</hbase.version>
</properties>
>
<!-- Flink -->
>
>org.apache.flink</groupId>
>flink-java</artifactId>
>${flink.version}</version>
</dependency>
>
>org.apache.flink</groupId>
>flink-streaming-java_${scala.binary.version}</artifactId>
>${flink.version}</version>
</dependency>
>
>org.apache.flink</groupId>
>flink-clients_${scala.binary.version}</artifactId>
>${flink.version}</version>
</dependency>
<!-- HBase -->
>
>org.apache.flink</groupId>
>flink-connector-hbase-2.2_${scala.binary.version}</artifactId>
>${flink.version}</version>
</dependency