简介:这是一套面向大数据开发初学者与进阶工程师的Flink实战项目源码,基于保险行业真实业务场景,采用Flink+HBase+Kafka+Phoenix架构,实现业务系统数据库数据的实时同步与实时统计报表分析,适合想通过完整项目理解流处理链路、补齐工程经验的学习者。资源包共73个文件,约584KB,以34个Java源码为核心,辅以11个Python脚本、9个Shell脚本,以及CSV数据样例、properties与XML配置、SQL建表语句、Phoenix初始化文件、Jar依赖和说明文档,覆盖从数据采集、消息队列到存储查询的完整环节。目前已有1166人学习下载。通过该资源可掌握实时同步任务的编码组织方式、Kafka与Flink的对接配置、HBase与Phoenix的存储查询设计,并借助脚本与配置快速搭建本地运行环境,理解真实项目中模块划分与排错思路。
1. 保险实时数仓为什么选 Flink:从 klrtdw 这个包说起
保险行业的实时数据链路有个很现实的特点:业务库分散、表结构经常改、下游报表口径又天天变。我拿到klrtdw.zip这个包的时候,第一反应不是看代码,而是先确认它到底解决了哪一段。解压后目录很干净:src/main下分prod和test两套资源,resources里放配置,根目录一个pom.xml、一个klrtdw.iml、一个test.log,外加README.md。这不是一个玩具 demo,而是一个按生产/测试双环境组织的 Maven 工程。
它要干的事很明确:用 Flink 把业务系统数据库的变更实时同步出来,落到 HBase,再通过 Phoenix 提供 SQL 查询能力,中间用 Kafka 做缓冲和解耦,最终支撑实时统计报表。技术栈是 Flink + Kafka + HBase + Phoenix 这套组合。适合谁?正在做 Flink 实时计算入门、想找一个带真实业务背景(保险)练手的同学,以及需要理解 CDC 同步 + 实时报表这条完整链路的中级开发。如果你只写过 WordCount,这个包能让你看到生产工程长什么样。
2. 环境搭建与工程结构:把 klrtdw 跑起来的第一步
2.1 从 pom.xml 反推依赖版本与组件选型
不要急着mvn clean package,先读pom.xml。这个文件决定了你本地要装什么版本的 Flink、Kafka 客户端、HBase 客户端。保险类项目通常对版本一致性极其敏感,Flink 的 minor 版本差异就可能导致连接器行为不同。
<!-- 典型依赖结构,具体版本以你包内 pom.xml 为准 --> <properties> <flink.version>1.13.x</flink.version> <!-- 以实际为准 --> <scala.binary.version>2.11</scala.binary.version> </properties> <dependencies> <!-- Flink 核心与流处理 --> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-streaming-java_${scala.binary.version}</artifactId> <version>${flink.version}</version> <scope>provided</scope> <!-- 集群已有,打包时不带入 --> </dependency> <!-- Kafka 连接器 --> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-connector-kafka_${scala.binary.version}</artifactId> <version>${flink.version}</version> </dependency> <!-- HBase / Phoenix 客户端 --> <dependency> <groupId>org.apache.phoenix</groupId> <artifactId>phoenix-core</artifactId> <version>5.x.x</version> <!-- 必须与集群 HBase 版本匹配 --> </dependency> </dependencies>逻辑说明:provided作用域的 Flink 核心包意味着你打出来的 jar 不含 Flink 本身,提交到集群时由集群提供,这是生产标准做法。参数上最需要盯的是三处——Flink 版本、Scala 二进制版本、Phoenix 版本。Phoenix 和 HBase 的版本必须严格对应,差一个小版本就可能出现NoSuchMethodError,这是血泪经验。
2.2 prod 与 test 双资源配置的切换逻辑
src/main/resources下分prod和test,这是 Maven profile 的经典用法。你需要确认pom.xml里有没有对应的<profiles>段,以及打包时激活哪个。
# 打包测试环境配置 mvn clean package -Ptest -DskipTests # 打包生产环境配置 mvn clean package -Pprod -DskipTests逻辑说明:-P激活 profile,Maven 会把对应目录下的配置文件打进 jar。参数上注意-DskipTests只是跳过测试执行,不是跳过编译。如果你不确定 profile 名,直接mvn help:active-profiles看当前激活了哪些。常见翻车点是两个环境的 Kafka broker 地址写反,导致本地测试连到生产,这个后面避坑章节细说。
2.3 本地跑通的最小依赖清单
在提交集群之前,建议本地先跑通。你需要:JDK 8(Flink 1.13 及以前基本锁 8)、Maven 3.6+、一个本地或可访问的 Kafka、一个 HBase 单机或伪分布式。HBase 单机版启动后默认端口 2181(ZooKeeper)和 16000(Master)。
# 启动本地 HBase(伪分布式) start-hbase.sh # 验证 HBase 可用 hbase shell > status > list逻辑说明:status返回集群节点数,list列出所有表。如果list卡住,八成是 ZooKeeper 没起来或hbase-site.xml里hbase.rootdir指向了不存在的 HDFS 路径。参数上,本地测试可以把hbase.rootdir设成file:///tmp/hbase避免依赖 HDFS。
3. 实时同步链路拆解:Kafka 到 HBase 再到 Phoenix
3.1 数据从业务库到 Kafka 的接入方式
保险业务库的变更要进 Kafka,常见做法有两种:一是用 CDC 工具直接抓 binlog 投递到 Kafka,二是业务层双写。这个项目走的是前者思路,Kafka 里存的是一条条变更事件。你需要先确认 topic 命名和分区数。
# 创建业务变更 topic,3 分区 1 副本(本地测试) kafka-topics.sh --create \ --bootstrap-server localhost:9092 \ --topic insurance_cdc \ --partitions 3 \ --replication-factor 1 # 查看 topic 详情 kafka-topics.sh --describe \ --bootstrap-server localhost:9092 \ --topic insurance_cdc逻辑说明:分区数决定 Flink 侧最大并行度,副本数在生产至少 2。参数上insurance_cdc只是示例名,实际以你resources里的配置为准。常见做法是按业务表拆多个 topic,避免单 topic 过大导致消费倾斜。
3.2 Flink 消费 Kafka 并写入 HBase 的核心算子
这是整个项目的骨架。Flink 从 Kafka 读流,做转换后通过 HBase 客户端写入。注意 HBase 写入不是靠 Flink 官方连接器,而是自己实现RichSinkFunction。
public class HBaseSink extends RichSinkFunction<String> { private transient Connection connection; @Override public void open(Configuration parameters) throws Exception { // 每个并行子任务独立建立连接,避免共享连接线程安全问题 org.apache.hadoop.conf.Configuration conf = HBaseConfiguration.create(); conf.set("hbase.zookeeper.quorum", "localhost"); conf.set("hbase.zookeeper.property.clientPort", "2181"); connection = ConnectionFactory.createConnection(conf); } @Override public void invoke(String value, Context context) throws Exception { // value 为解析后的行键与列值,按业务拼装 Put Table table = connection.getTable(TableName.valueOf("ins_policy")); Put put = new Put(Bytes.toBytes(extractRowKey(value))); put.addColumn(Bytes.toBytes("info"), Bytes.toBytes("amount"), Bytes.toBytes(extractAmount(value))); table.put(put); table.close(); } @Override public void close() throws Exception { if (connection != null) connection.close(); } }逻辑说明:open里建连接、close里释放,这是RichSinkFunction的标准生命周期。参数上hbase.zookeeper.quorum必须和集群一致。关键坑在于Table对象不要每次invoke都新建又关闭,高频写入下这会拖垮性能,正确做法是复用 Table 或在open里初始化。上面为了展示清晰才在invoke里取表,实际应提到open。
3.3 Phoenix 建表与实时报表查询
HBase 原生 API 查询不友好,Phoenix 在 HBase 上套了一层 SQL。建表时要注意 Phoenix 的 rowkey 设计直接决定查询性能。
-- Phoenix 建表,rowkey 为保单号 CREATE TABLE IF NOT EXISTS ins_policy ( policy_no VARCHAR NOT NULL PRIMARY KEY, amount DECIMAL(18,2), create_time TIMESTAMP ); -- 实时统计:按小时汇总保费 SELECT TO_CHAR(create_time, 'yyyy-MM-dd HH') AS hour_bucket, SUM(amount) AS total_amount FROM ins_policy GROUP BY TO_CHAR(create_time, 'yyyy-MM-dd HH');逻辑说明:Phoenix 的PRIMARY KEY就是 HBase rowkey,GROUP BY会走全表扫描,数据量大时很慢。参数上DECIMAL(18,2)对应金额精度。常见做法是把时间维度做进 rowkey 前缀(如policy_no反转或加盐),避免热点写入。报表查询如果频繁按时间聚合,建议预聚合或建二级索引。
4. 避坑与排查:这几个问题我踩过不止一次
4.1 现象:任务提交后报 ClassNotFoundException
原因:Flink 核心依赖用了provided,但某些连接器(如 Kafka)没打进去,或者打包时漏了 shade 插件。解决:确认pom.xml里非 provided 的依赖都被打进 jar,用mvn dependency:tree排查冲突,必要时加maven-shade-plugin并指定Main-Class。
4.2 现象:HBase 写入报 RegionTooBusyException
原因:rowkey 设计单调递增(如直接用时间戳或自增 ID),所有写入压到同一个 Region。解决:rowkey 加盐或反转,比如policy_no前加两位哈希前缀,把写入打散到多个 Region。
4.3 现象:Kafka 消费延迟越来越高,checkpoint 频繁失败
原因:HBase 写入是同步阻塞的,invoke里每次建连接或 Table 导致吞吐上不去,反压传导到 Kafka。解决:连接和 Table 在open里初始化并复用,写入改为批量BufferedMutator,同时调大 checkpoint 超时。
4.4 现象:Phoenix 查询报 TableNotFoundException 但 HBase 里明明有表
原因:Phoenix 的 schema 映射和 HBase 原生表不一致,Phoenix 建的表在 HBase 里是带特殊字符的。解决:统一用 Phoenix 建表,不要用 HBase shell 建了再让 Phoenix 去读,两者元数据不通用。
4.5 现象:本地跑通,集群上配置文件读的是 prod 但连的是 test 地址
原因:profile 激活错误或resources目录下同名文件覆盖顺序问题。解决:打包后unzip -l target/xxx.jar | grep resources确认实际打进去的是哪套配置,别凭感觉。
5. 进阶技巧:用 test.log 反推数据流与验证方法
test.log这个文件很多人会忽略,但它其实是验证链路是否跑通的最快入口。我一般会先tail -f test.log,观察 Flink 任务的启动日志、Kafka 消费 offset、HBase 写入批次。如果日志里出现连续的invoke耗时打印,说明写入是瓶颈;如果 offset 长时间不动,说明上游没数据或反压了。
验证方法上,我习惯用「三段对账」:Kafka 里kafka-run-class.sh kafka.tools.GetOffsetShell看某 topic 总消息数,Flink 的 Web UI 看numRecordsIn和numRecordsOut,Phoenix 里SELECT COUNT(*)看落库条数。三个数对不上,就按链路逐段排查。参数上注意 Flink UI 的并行度和 Kafka 分区数的关系,并行度大于分区数时会有子任务空转。
# 查看 Kafka topic 各分区最新 offset kafka-run-class.sh kafka.tools.GetOffsetShell \ --broker-list localhost:9092 \ --topic insurance_cdc --time -1 # Phoenix 侧核对落库总量 sqlline.py localhost:2181:/hbase > SELECT COUNT(*) FROM ins_policy;逻辑说明:--time -1取最新 offset,-2取最早。Phoenix 的sqlline.py连接串格式是zk:port:/hbase。这两个数加上 Flink UI 的计数,基本能定位问题出在消费、转换还是写入环节。
从那以后我每次拿到一个新的 Flink 工程,都强制先跑一遍test.log观察 + 三段对账,再动业务代码。这个习惯帮我省了太多「代码看着没问题但数据就是不对」的排查时间。希望帮到你。
本文还有配套的精品资源,点击获取