ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

保险行业实时报表:Kafka到HBase的Spark流处理全链路实践

保险行业实时报表:Kafka到HBase的Spark流处理全链路实践 简介面向大数据开发与Spark学习者的保险行业实时数据分析实战项目包基于Kafka、Spark Streaming与HBase构建完整数据管道演示数据库变更实时同步、实时聚合统计与报表查询等典型场景适合希望掌握实时计算全链路开发的初中级工程师。压缩包共400个文件约632KB以28个Scala源码、9个Shell脚本、8个Properties配置、7个XML文件为主另有11个sample样例文件及Git仓库内部文件HEAD、master、config等便于查看工程版本演进脉络。目前已有981人学习项目涵盖Kafka生产者与消费者、Spark流式窗口计算、HBase结果持久化与查询接口等模块可帮助读者理解各组件如何协同工作并直接参考其工程目录结构与启动脚本。压缩包体积虽小但代码配置密度高适合边读边实践快速搭建真实保险业务场景的实验环境稍加修改即可将主题与表结构替换为其他业务场景复用于实时报表与指标监控。1. 保险行业真实 Spark 项目从 Kafka 到 HBase一条实时报表流水线的全部细节保险公司每天早晨跑批产生的 T1 报表放到实时决策场景里经常是“马后炮”车险报案量激增、健康险理赔异常走高等到第二天看到报表已经错过干预窗口。我拆过的这个保险行业真实项目核心就是把业务系统数据库的变更数据实时采集到 Kafka由 Spark Streaming 完成清洗、窗口聚合结果写入 HBase再对接报表接口。链路听起来不复杂但从集群版本匹配到 RowKey 设计处处都能让人翻车。下面把整条流水线拆开从架构选型讲到参数踩坑适合正在做实时数仓或准备把流处理引入保险业务场景的同学。2. 整体架构与技术选型为什么是 Kafka Spark HBase 三件套实时报表的链路里Kafka、Spark、HBase 的组合出现频率极高不是因为它最“新”而是因为生态最成熟。保险行业的数据链路往往还要兼容 Oracle、MySQL 等异构数据源Canal 或 Debezium 订阅 binlog 后统一发给 KafkaSpark 从 Kafka 消费处理后写 HBase。选型的核心判断标准有两个第一消息中间件要能扛住业务库变更的瞬时高峰第二存储层要满足随机查询和范围扫描。下面把这套组合拆开讲。2.1 业务指标拆解保险实时报表到底要算什么保险实时报表不能套用互联网通用的“日活”“GMV”那套逻辑。一个保险核心系统里实时需要看的指标通常分成三类保费收入按渠道、险种、时间段统计累计保费。理赔进展报案量、结案量、理赔金额、件均赔款。核保动态自动核保通过率、人工核保积压量。其中保费收入是“当日累计”口径理赔进展是“最近一小时窗口”口径核保动态更接近“实时明细查询”。不同的口径对计算模型的要求完全不一样。我一般会把指标切成两组一组走“明细累计”直接用 HBase 的 increment 或 update一组走“窗口聚合”用 Spark Streaming 的 window 操作。这样做的目的是把计算模型简化能累加的尽量不进窗口能窗口聚合的不要做全局去重。如果你一开始就把所有指标揉进同一个窗口计算后面数据量上来再加维度改动成本会成倍增加。2.2 Kafka 的定位与核心参数削峰填谷和消息不丢Kafka 在这条链路里不是简单的“中转站”它同时承担了削峰填谷和故障缓冲两件事。业务系统在白天高峰期的数据库事务量和凌晨跑批完全不同如果没有 KafkaSpark 要按峰值设计资源平时就会严重浪费。Kafka 把消息持久化在磁盘上Spark 按自己的处理能力消费等于在两个系统之间加了一层弹性缓冲。生产环境我通常会重点校对这些参数参数推荐值作用replication.factor3允许一台 broker 故障不丢数据min.insync.replicas2配合 acksall 保证至少两个副本写入成功acksall生产者侧保证消息提交成功retention.ms17280000048 小时给报表任务留足重跑窗口compression.typelz4减小网络带宽占用关于 retention 多说一句很多同学默认 Kafka 保留 7 天其实在实时报表链路里数据一旦被 Spark 消费并写入 HBaseKafka 里的原始消息就没太大用了。保留 48 小时是为了解决“凌晨发现昨天数据算错了想要重跑”的需求再长只会占磁盘。如果你有离线分析的需求可以用另一个 topic 单独长期保存不要让实时 topic 背上存储包袱。2.3 Spark 流处理选型DStream 和 Structured Streaming 怎么选Spark 侧最容易纠结的是用老的 DStream API 还是新的 Structured Streaming。这个项目的第一个版本用 DStream后来重构时换成了 Structured Streaming原因有三点。Structured Streaming 把“流”抽象成一张不断追加的表查询逻辑和 Spark SQL 几乎一样写窗口聚合、join、过滤都比 RDD 时代的 DStream 直观不少。第二它内置了端到端的 exactly-once 语义只要 sink 侧配合自己做幂等就能做到“不重不丢”。第三新版本的foreachBatch允许复用 DataFrame 的批处理能力对一个微批次做 collect、persist、跨服务写入非常方便。不过 DStream 也不是完全过时。如果你的团队对 RDD 和foreachRDD已经很熟而且任务逻辑简单只有“消费-聚合-写入”三步DStream 的维护成本可能更低。选择的关键依据是团队的技能树而不是“哪个新用哪个”。生产上我见过不少 DStream 任务跑得远比 Structured Streaming 稳的例子换了 API 但业务模型没变反而暴露了一堆隐藏问题比如 checkpoint 格式不兼容和变更后需要重构整个算子流。2.4 HBase 存储选型为什么不用 MySQL 或 Redis实时报表的结果数据有两个特点一是写入量巨大每个窗口每个维度都会产生一条数据二是查询集中通常按渠道和时间范围批量扫描。MySQL 在亿级行数面前容易因为索引和磁盘 IO 成为瓶颈Redis 的内存容量和大数据量成本又太高HBase 的分布式列存储正好匹配这种“大规模、稀疏、随机读写”的场景。RowKey 的设计直接决定了 HBase 写入和查询的性能。保险报表里我常用的 RowKey 格式如下{指标分类}_{业务日期}_{渠道ID}_{窗口开始时间}例如PREMIUM_20251224_AGENT01_202512241500表示 2025 年 12 月 24 日 15:00 这个窗口、AGENT01 渠道的保费累计。这个设计有两个好处同一渠道同一时间段的数据在物理上相邻Scan 效率高渠道 ID 不是强递增序列写入时可以分散到不同 Region。要注意的是如果渠道 ID 本身有“大渠道效应”比如某个线上渠道的数据量占 70%那按渠道前缀的 RowKey 就会产生热点写。这个问题的解法后面第 5 章会专门讲。3. 环境搭建集群规划、版本匹配与三类组件的关键配置先声明生产级集群不会只有一台机器但学习和验证阶段最少需要 3 台节点。我提供的这套配置是我自己落地过的组合Spark 3.2.2 Kafka 2.8.1 HBase 2.4.9 Hadoop 3.3.2。这套组合的主要考量是避开 Spark 3.2 和 HBase 2.2 之间的一些 API 兼容问题同时 Spark 3.2.2 已经是 3.2 线最稳定的 Patch 版本。3.1 版本匹配先把最容易被忽略的坑抠掉大数据组件之间的版本兼容一直是最容易踩坑的地方。我之前在一次项目里用 Spark 3.0 配 HBase 2.1结果hbase-spark模块因为 Hadoop 版本差异直接编译不过额外耗了两天。后来总结出下面这套相对稳妥的版本组合组件版本用途Hadoop3.3.2HDFS 做 HBase 持久化和 Spark shuffle 存储HBase2.4.9报表数据存储Java API 读取Kafka2.8.1业务 binlog 消息管道Spark3.2.2流处理引擎Scala 2.12Zookeeper3.6.3Kafka 和 HBase 共用协调服务这里要特别注意一点HBase 的hbase-spark模块已经很久没有更新新项目不太建议依赖它做 DataFrame 映射直接用 HBase 客户端 API 自己写读写反而更可控。Spark 3.2.2 和 Kafka 0-10 连接器不需要额外引入 spark-streaming-kafka只要在依赖里加上spark-sql-kafka-0-10_2.12:3.2.2即可。3.2 Kafka Topic 初始化分区数、副本数和 retention 的设定Kafka broker 部署好之后先建 topic。下面的命令创建一个 6 分区、3 副本的 topic专门承载保险业务库保单表的 binlog 变更消息。kafka-topics.sh --bootstrap-server kafka01:9092,kafka02:9092,kafka03:9092 \ --create \ --topic insurance_biz_policy \ --partitions 6 \ --replication-factor 3 \ --config retention.ms172800000 \ --config min.insync.replicas2参数说明--partitions 6不是拍脑袋定的。后续 Spark 实时任务的 executor 总核心数如果是 12分区数 6 ~ 12 之间都能获得不错的消费并行度。--replication-factor 3要求集群有至少 3 个 broker副本数大于 3 就是浪费磁盘。retention.ms172800000表示消息保留 48 小时这个窗口要覆盖“发现数据异常并重跑当天任务”的时间。Topic 创建完成后可以用kafka-console-consumer.sh快速验证消息能否正常生产消费验证 binlog 采集端是否已经打通kafka-console-consumer.sh --bootstrap-server kafka01:9092 \ --topic insurance_biz_policy \ --from-beginning \ --max-messages 5如果这里能拉到 Canal 或 Debezium 生产的 JSON 消息说明采集链路是通的如果一直阻塞优先排查 Kafka 和 Canal 所在的服务器网络、advertised.listeners配置是否正确。3.3 HBase 预分区建表列族、TTL 与 Region 分布HBase 生产环境建表最忌讳直接使用默认配置。默认创建的表只有一个 Region所有写入都打到一台 RegionServer数据量和吞吐量稍微上来就直接热点写崩。下面的建表语句把insurance_report预分为 10 个 Region同时设置了两个列族。hbase shell EOF create insurance_report, {NAME cf_metric, VERSIONS 1, TTL 604800, COMPRESSION SNAPPY}, {NAME cf_meta, VERSIONS 1, TTL 604800, COMPRESSION SNAPPY}, {NUMREGIONS 10, SPLITALGO HexStringSplit} EOF逻辑说明cf_metric存实际指标值cf_meta存数据版本、写入时间等元信息。两个列族都设了 7 天 TTL过期数据自动清理防止报表表无限膨胀。NUMREGIONS 10在 HBase 0.94 之后的 shell 里可以直接生效底层通过SplitAlgorithm把 RowKey 空间切割成 10 段。这里有一个前置条件只有当 RowKey 的前缀能够均匀分布在 0x00-0xff 区间时HexStringSplit 才有意义。如果你用的是类似于“日期渠道”的 RowKey前面那位是业务日期数字散列效果并不好。所以我通常建议在 RowKey 最前面加一个 0~9 的业务分桶位配合SPLITALGO UniformSplit更合理。分桶位从哪来最简单的做法是用channel_id.hashCode() % 10的绝对值。3.4 Spark 提交参数既要吞吐也要稳定性实时任务在 YARN 上提交时executor 资源、Kafka 消费速率、shuffle 并行度这几个参数会直接影响任务的生死。下面是我在这套项目里稳定跑了几周的提交脚本spark-submit \ --master yarn \ --deploy-mode cluster \ --name insurance_realtime_report \ --executor-memory 8g \ --executor-cores 4 \ --num-executors 3 \ --conf spark.streaming.kafka.maxRatePerPartition1000 \ --conf spark.streaming.backpressure.enabledtrue \ --conf spark.sql.shuffle.partitions24 \ --conf spark.serializerorg.apache.spark.serializer.KryoSerializer \ --conf spark.kryoserializer.buffer.max512m \ --conf spark.shuffle.service.enabledtrue \ --jars hbase-client-2.4.9.jar,hbase-common-2.4.9.jar,hbase-server-2.4.9.jar \ insurance-realtime.jar这段是实时任务最容易被抄错的地方。maxRatePerPartition1000限制每个分区每秒最多消费 1000 条当 Kafka 积压几十万条消息时任务不会一启动就被打爆。backpressure.enabledtrue让 Spark 根据上一批处理耗时动态调节消费速率两个参数加在一起等于给消费端上了“限流油门”。spark.sql.shuffle.partitions24是我用手头 3 个 executor、每个 4 核算出来的。默认值 200 在流处理任务里会产生大量小任务Shuffle 文件碎片化严重写 HBase 时也会因为 partition 过多造成小文件24 这个数字接近 executor 总核心数 12 的 2 倍聚合和写入阶段的并发恰到好处。如果你的数据倾斜比较明显可以继续往上调但要同步观察 HBase RegionServer 的 CPU 和 GC 指标。4. 核心实现消费、聚合、写入 HBase 的完整代码拆解环境准备好之后核心逻辑主要分三步读 Kafka、窗口聚合、写 HBase。下面用 Structured Streaming 的实现来说明代码以 PySpark 为准实际生产里逻辑相同。4.1 读取 Kafka 与 JSON Schema 解析Structured Streaming 读 Kafka 拿到的原始 DataFrame 里value 是二进制字节需要结合业务消息格式做解析。下面的代码定义了保险业务消息的 Schema并完成时间字段的转换。from pyspark.sql import SparkSession from pyspark.sql.functions import from_json, col, to_timestamp from pyspark.sql.types import (StructType, StructField, StringType, LongType, DoubleType) # 与 Canal 订阅表中的字段保持严格一致 biz_schema StructType([ StructField(table_name, StringType()), StructField(op_type, StringType()), StructField(occur_time, LongType()), StructField(policy_id, StringType()), StructField(channel_id, StringType()), StructField(insurance_type, StringType()), StructField(premium_amount, DoubleType()), StructField(claim_amount, DoubleType()), StructField(report_status, StringType()) ]) spark SparkSession.builder \ .appName(insurance_realtime_report) \ .config(spark.sql.shuffle.partitions, 24) \ .config(spark.serializer, org.apache.spark.serializer.KryoSerializer) \ .getOrCreate() df_raw spark.readStream \ .format(kafka) \ .option(kafka.bootstrap.servers, kafka01:9092,kafka02:9092,kafka03:9092) \ .option(subscribe, insurance_biz_policy) \ .option(startingOffsets, earliest) \ .option(maxOffsetsPerTrigger, 10000) \ .option(failOnDataLoss, false) \ .load() df_biz df_raw.selectExpr( CAST(key AS STRING) AS biz_key, CAST(value AS STRING) AS biz_value ).select(from_json(col(biz_value), biz_schema).alias(data)) \ .select(data.*) \ .withColumn(biz_time, to_timestamp(col(occur_time) / 1000))逻辑说明from_json按 biz_schema 解析 Kafka value比让 Spark 自己推断 schema 稳定得多。occur_time是 Canal 采集到的 binlog 毫秒时间戳除以 1000 转成秒再传给to_timestamp以后窗口聚合直接使用biz_time字段。key字段可用于后续去重但这里没有直接参与计算。参数说明startingOffsetsearliest表示任务首次启动时从 topic 最早可消费位置开始适合初始化报表和历史数据补算如果只想从当前时刻开始改为latest。maxOffsetsPerTrigger10000限制每个批次最大消费 1 万条避免启动时瞬间拉取巨量数据导致后续聚合卡顿。failOnDataLossfalse表示允许 Kafka offset 因为过期被清理后不直接失败配合监控告警使用不要盲目依赖它掩盖问题。4.2 窗口聚合与多维度指标计算保险报表的核心指标是保费累计和理赔进展。下面的代码用 5 分钟滚动窗口按渠道和险种维度聚合premium_sum、claim_sum、报案量和结案量。from pyspark.sql.functions import window, sum, count, when df_agg df_biz \ .filter(col(op_type) ! DELETE) \ .groupBy( window(col(biz_time), 5 minutes), col(channel_id), col(insurance_type) ) \ .agg( sum(premium_amount).alias(premium_sum), sum(claim_amount).alias(claim_sum), count(when(col(report_status) REPORTED, 1)).alias(report_cnt), count(when(col(report_status) SETTLED, 1)).alias(settle_cnt) ) \ .withColumn(window_start, col(window.start)) \ .withColumn(window_end, col(window.end))逻辑说明过滤DELETE操作是为了避免业务数据删除导致的负向金额。window(biz_time, 5 minutes)生成滚动窗口每个窗口输出window_start和window_end下游写 HBase 时可以使用window_start拼 RowKey。条件计数count(when(...))对满足特定状态的消息计数一条消息只会累加到对应状态字段。这里有几个容易踩的细节窗口聚合默认会在窗口结束后才输出延迟取决于 watermark 配置实时报表可以接受 1~2 分钟延迟我就用withWatermark(biz_time, 2 minutes)来控制乱序数据的清理op_type如果是 UPSERT 语义需要额外处理更新前的旧值这里按当前值直接覆盖如果业务要求“当日累计”而不是窗口累计就不能用 window 字段做 groupBy 的顶层应该单独按日期 渠道做 groupBy用批内增量累加到 HBase 已有值上。4.3 结果写入 HBaseforeachBatch 与批量 PutStructured Streaming 没有官方的 HBase sink生产上最通用的姿势是用foreachBatch把每个微批 DataFrame 收集后通过 HBase Java API 写入。下面的代码展示了核心写入逻辑。from hbase_utils import HBaseWriter def write_metric_to_hbase(df_epoch, epoch_id): 每个微批回调一次epoch_id 可用于日志追踪和去重 df_epoch.persist() df_epoch.localCheckpoint() rows df_epoch.collect() if not rows: df_epoch.unpersist() return puts [] for row in rows: row_key PREMIUM_{}_{}_{}.format( row[window_start].strftime(%Y%m%d%H%M), row[channel_id], row[insurance_type] ) puts.append(HBaseWriter.build_put( row_key, {cf_metric:premium_sum: str(row[premium_sum]), cf_metric:claim_sum: str(row[claim_sum]), cf_metric:report_cnt: str(row[report_cnt]), cf_metric:settle_cnt: str(row[settle_cnt])} )) HBaseWriter.batch_put(insurance_report, puts) df_epoch.unpersist() df_agg.writeStream \ .foreachBatch(write_metric_to_hbase) \ .outputMode(update) \ .trigger(processingTime30 seconds) \ .start() \ .awaitTermination()逻辑说明foreachBatch在每个 trigger 结束后把本次输出的 DataFrame 交给回调函数。localCheckpoint()是防止失败恢复后同一个批次数据重算再写入 HBase 的关键它把当前批的中间状态保存到本地磁盘下次同一 batch 会直接取 checkpoint 的结果。数据行通过 HBaseWriter 的build_put转成 Put 对象再batch_put一次性提交比单条写入快一个数量级。参数说明outputMode(update)表示每个窗口的最新聚合结果会随新到的数据不断更新配合 HBase 按 RowKey 覆盖写入就能实现“近实时更新”。trigger(processingTime30 seconds)是触发间隔如果业务可以接受秒级延迟就不需要调小Kafka 分区多、数据量大的时候把 trigger 调小到 5 秒也是常见做法但要保证 executor 能在触发间隔内处理完上一批数据否则任务会堆积。4.4 HBase 查询层的 RowKey 路由与 Scan 过滤报表服务读 HBase 时最常用的查询是按渠道和日期范围查当天所有时段的指标。这里的 RowKey 前缀可以直接传给 ScanScan scan new Scan(); scan.withStartRow(Bytes.toBytes(PREMIUM_20251224_AGENT01_)); scan.withStopRow(Bytes.toBytes(PREMIUM_20251224_AGENT01_ \u0000));逻辑说明withStartRow包含前缀withStopRow再加上\u0000作为后缀截断可以扫出该前缀下的全部分区。如果查询维度不含渠道那就只能扫描整张报表表性能会差很多。所以在查询需求不明确之前先把“必查维度”都拆出来作为 RowKey 前缀的一部分。5. 避坑指南Kafka 消费、HBase 写入和 Spark 调优的 5 个血泪教训再简单的实时任务上了生产都会变得很复杂。下面 5 个问题是我在保险项目上线前后真实踩过、并且排了一整天才定位的坑按“现象 → 原因 → 解决”的顺序写方便排查时直接对号入座。5.1 任务重启后重复消费HBase 里指标翻倍现象某天凌晨 YARN 资源紧张导致任务被 kill管理员把它拉起来后第二天对账发现保费收入比业务系统多了 11%。原因Structured Streaming 默认从 checkpoint 恢复 offset但foreachBatch写 HBase 发生在 Spark 提交 offset 之前。任务 fail 在“写完 HBase”和“提交 offset”之间时同一批消息会被再次消费并再次写入 HBase导致重复数据。解决写入 HBase 前做幂等。我对每条 RowKey 加窗口时间和渠道险种写入前用 HBase 的checkAndMutate判断如果该 RowKey 已存在且新数据的时间戳更大才覆盖。注意checkAndMutate的语义是“当列值满足某个条件时才执行修改”可以把数据版本号放到cf_meta:version列比对当前版本是否大于已存版本用原子操作避免并发写入时重复累加。补充一个排查方法当发现数据翻倍时先看 Spark UI 里这个 batch 有没有重试记录再看 HBase 里同一个 RowKey 是否有两条时间戳不同的数据。通过HBase shell count和scan很快就能确认是哪批数据重复进来。5.2 Spark 任务频繁 OOMGC 日志全是 Full GC现象任务稳定运行 2 小时后executor 日志报java.lang.OutOfMemoryError: Java heap space数据量远没到集群上限但 GC 日志里全是 Full GC。原因spark.sql.shuffle.partitions用了默认的 200聚合结果的中间文件被散到 200 个分区每个分区只处理很少的数据产生了大量小任务executor 上频繁创建和销毁对象让 JVM 压力骤增。还有一个隐蔽原因是collect()之前没有对 DataFrame 做持久化每次 collect 都会重建上游全部聚合结果导致重复 GC。解决根据集群规模调整spark.sql.shuffle.partitions我的经验值是 executor 总核心数的 1.5~2 倍。同时给 shuffle read 阶段预留足够内存配置spark.shuffle.service.enabledtrue把溢写数据交给独立 shuffle 服务减轻 executor 反复拉取数据的负担。另外在 foreachBatch 里一定要persist()collect()unpersist()保证同一微批内最多执行一次完整的血缘计算。5.3 HBase Region 热点一个 RegionServer 被写爆其他很闲现象监控面板上某个 RegionServer 的写请求 QPS 是其他节点的 8 倍磁盘 IO 打满整个 HBase 集群都感觉变慢甚至出现 RegionServer 挂掉后自动重启。原因RowKey 以渠道 ID 开头互联网大渠道的数据量远大于传统渠道所有数据大概率落到同一段 RowKey 区间跑到了同一个 Region 上。热点写不单影响该 RegionServer还会触发 memstore 频繁 flush产生大量小 HFile进一步拖慢查询。解决RowKey 前增加散列前缀例如hash(channel_id) % 9作为首字符让不同渠道均匀分布在多个 Region。查询时通过前缀路由表定位到具体的前缀再拼接渠道和日期扫描。这个改动对报表查询接口的影响不大但能显著缓解写入热点。要注意散列前缀的位数和预分区数目对齐预分区 10 个前缀范围就取 0~9否则会再次出现“部分 Region 无数据、部分 Region 被打满”的情况。5.4 报表数据和业务库对不上少了一个零头现象某天增速报表里保费收入比业务库的实时 sum 少了 0.3%查了半天没发现代码逻辑有明显 bug。原因业务库对一条保单做了 DELETECanal 把 DELETE 操作也推到了 Kafka。Spark 任务没有过滤 DELETE直接按负值累加进保费导致报表少计了金额。另外Canal 在启停时还会产生重复的 EVENT 标记如果采集线程发生过网络重连同一条 binlog 可能被推送两次。解决解析阶段显式过滤op_type DELETE同时在写入 HBase 前把premium_amount和claim_amount的 null 值填充为 0。对重复消费问题要在生产链路里持续观察occur_time和 canal 内部的executeTime如果同一条 policy_id 在同一个窗口内出现两次后到的版本更新前一条即可不要简单累加。5.5 checkpoint 跨集群迁移新集群任务反复失败现象把任务从测试集群迁移到生产集群Spark UI 显示任务提交链路没问题但持续报FileNotFoundException指向 checkpoint 目录里的某些二进制文件。原因checkpoint 里保存了原集群的 HDFS 路径、Executor 元数据和部分序列化算子信息。直接把旧的 checkpoint 目录复制到新集群任务恢复时找不到原集群的中间文件自然失败。解决迁移集群时删除 checkpoint并显式指定startingOffsets。如果不想从头消费 Kafka可以从旧集群导出每个 topic 分区的 offset然后通过startingOffsets选项精确指定新任务的消费起始位置。不要把 checkpoint 当成备份来复制这个文件目录只对“同一集群下的同一任务”有意义。6. 进阶调优与验证从“能跑”到“跑得稳”的最后一公里实时任务上线前和后期的稳定性维护远比“把窗口计算写出来”重要。下面这块是我养成的习惯也解决过多次凌晨的惊醒。6.1 上线前必做全链路对账实时报表上线前我会先把 Kafka 生产端、Spark 消费端、HBase 总量三个数字对清楚。做法很简单在某个窗口内手工给 Kafka topic 灌入一笔已知金额的测试数据然后在 HBase 查出对应 RowKey 的值手工加总与预期核对。对账表大致长这样检查项工具通过标准Kafka 生产总量kafka-consumer-groups.sh与写入的测试消息数量一致Spark 消费总量Spark Streaming 指标累计 input 数等于 Kafka 消息数HBase 指标和HBase Shell scan对账金额与测试值一致报表接口返回直接调查询接口返回时间 200ms不要小看这一步很多“数据对不上”的问题只要提前做半小时对账就能拦住。6.2 监控与告警的配置习惯上线前我会配三类告警一是 Kafka lag 告警超过 5 万就通知说明消费速度跟不上生产速度二是 Spark Streaming 的 batch 处理耗时告警连续 3 个 batch 超过触发间隔就报警三是 HBase RegionServer 的写队列长度告警超过阈值说明 RowKey 或者预分区出了问题。这三类告警用 Grafana Prometheus 就能接起来Kafka exporter、Spark metrics、HBase 的 JMX 都有现成模板关键是阈值要根据数据量级别自己调不要用模板默认值。最后一个习惯是从那以后每次上线我都强制走一遍全链路对账 三天保留期的 Kafka retention 检查。保留期太短凌晨重跑就只能干瞪眼太长又浪费磁盘。磁盘余量打到 60% 以下的那一天就着手扩容 broker别等磁盘满告警再动手。希望帮到你。本文还有配套的精品资源点击获取
返回列表