
简介这是一套面向大数据开发初学者与进阶工程师的Flink实战项目源码基于保险行业真实业务场景采用FlinkHBaseKafkaPhoenix架构实现业务系统数据库数据的实时同步与实时统计报表分析适合想通过完整项目理解流处理链路、补齐工程经验的学习者。资源包共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.version1.13.x/flink.version !-- 以实际为准 -- scala.binary.version2.11/scala.binary.version /properties dependencies !-- Flink 核心与流处理 -- dependency groupIdorg.apache.flink/groupId artifactIdflink-streaming-java_${scala.binary.version}/artifactId version${flink.version}/version scopeprovided/scope !-- 集群已有打包时不带入 -- /dependency !-- Kafka 连接器 -- dependency groupIdorg.apache.flink/groupId artifactIdflink-connector-kafka_${scala.binary.version}/artifactId version${flink.version}/version /dependency !-- HBase / Phoenix 客户端 -- dependency groupIdorg.apache.phoenix/groupId artifactIdphoenix-core/artifactId version5.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激活 profileMaven 会把对应目录下的配置文件打进 jar。参数上注意-DskipTests只是跳过测试执行不是跳过编译。如果你不确定 profile 名直接mvn help:active-profiles看当前激活了哪些。常见翻车点是两个环境的 Kafka broker 地址写反导致本地测试连到生产这个后面避坑章节细说。2.3 本地跑通的最小依赖清单在提交集群之前建议本地先跑通。你需要JDK 8Flink 1.13 及以前基本锁 8、Maven 3.6、一个本地或可访问的 Kafka、一个 HBase 单机或伪分布式。HBase 单机版启动后默认端口 2181ZooKeeper和 16000Master。# 启动本地 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 再到 Phoenix3.1 数据从业务库到 Kafka 的接入方式保险业务库的变更要进 Kafka常见做法有两种一是用 CDC 工具直接抓 binlog 投递到 Kafka二是业务层双写。这个项目走的是前者思路Kafka 里存的是一条条变更事件。你需要先确认 topic 命名和分区数。# 创建业务变更 topic3 分区 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 RichSinkFunctionString { 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 rowkeyGROUP 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和numRecordsOutPhoenix 里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观察 三段对账再动业务代码。这个习惯帮我省了太多「代码看着没问题但数据就是不对」的排查时间。希望帮到你。本文还有配套的精品资源点击获取