ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

Kafka+Spark实时看板避坑:checkpoint、offset与幂等写入

Kafka+Spark实时看板避坑:checkpoint、offset与幂等写入 上个月有个做电商数据的朋友找我说他们花两周用 Spark Kafka 搭了一套实时 Dashboard上线第一天数字是对的第二天开始就乱了销售额比后台少了 8%第三天又多了 12%。他打开 Kafka 的可视化工具看了一眼lag 是 0消费组也在跑一切看起来都正常。这种情况我遇到过不止一次——Spark 和 Kafka 本身都没问题问题出在两者之间那条谁也没写进文档的缝里。这篇笔记就是把这套链路上真正会绊人的地方摊开讲清楚。从 Kafka 集群安装配置、Topic 分区怎么定到 Spark 内存调优、checkpoint 和 offset 的取舍再到 Dashboard 展示层的预聚合和刷新节奏我会把每一个选择的理由、每一个坑的排查链路都写出来。适合已经跑通过 Hello World 级别 demo、但把它放到真实数据量下就开始出问题的人也适合正准备用 Spark Kafka 做实时看板、想少走点弯路的同学。基础概念我会顺带补两句但不会花篇幅讲 Kafka 是什么、Spark 是什么这种教科书内容。1. 先把链路想清楚Kafka 到 Dashboard 之间要经过几次搬运1.1 直连数据库做看板卡在哪一步大多数人第一次做实时看板第一反应是让前端定时轮询业务库。这个方案在数据量小的时候能跑但它有三个天花板而且都是硬的。第一个是业务库压力。看板上的每一个数字背后都是一条聚合 SQL如果每 10 秒刷新一次、有五个人同时打开看板你的业务库上就挂着每秒几十条 GROUP BY。业务库的 CPU 是被下单、支付这些核心链路占着的看板把 CPU 抢走最先报警的一定是交易系统。第二个是明细计算的时间成本。假设你想看近 7 天复购率在明细表上算需要扫描七天的订单明细几百万行起步单次查询就得几秒。看板又要求秒级刷新这两件事天然矛盾。第三个是无法承载实时。业务库里的数据本身就是落库后的状态支付网关回调延迟、异步对账、状态机流转都会让业务库的数据比你想要的实时视图晚几分钟到几十分钟。有些指标比如曝光点击根本不会进业务库。所以引入 Kafka 的第一个价值是把事件流和业务库解耦事件先进 KafkaKafka 承担削峰和缓冲Spark 从 Kafka 里按自己的节奏消费算完之后写到一张专门给看板用的聚合表里。业务库不受影响看板也只查聚合表。1.2 批、微批、纯流三种形态怎么选Spark 处理 Kafka 数据具体有三种落地形态它们的取舍差别很大选错了后期改造成本很高。形态延迟量级实现成本故障恢复适用场景定时批Spark SQL Kafka 不参与分钟到小时最低简单重跑即可T1 报表、日报微批Structured Streaming 默认秒级中等依赖 checkpoint实时看板、监控告警纯流Continuous Processing毫秒级高生态受限极少一般不推荐我这里说的微批就是 Structured Streaming 默认的执行方式它本质上是一个调度器 小批次批处理每隔一个 trigger 间隔把 Kafka 里积累的新数据拉过来当成一个 DataFrame 处理。这个模型有个非常重要的好处计算逻辑可以复用批处理的写法你写好的 DataFrame 转换逻辑换成 readStream 就能跑调试成本极低。而 Continuous Processing 虽然延迟低但它对算子有严格限制聚合、join 很多都不支持而且失败恢复要重跑。我做过的项目里99% 的场景微批就够了——trigger 设成 10 秒或 30 秒业务上完全感知不到差别但开发效率和稳定性高一个数量级。1.3 指标口径先定再谈技术选型这一条是我踩过最贵的坑跟技术完全无关。朋友的复购率口径是这样定义的统计周期内产生第二笔及以上订单的用户数 ÷ 统计周期内下单的用户数。但运营那边的口径是统计周期内每个用户的下单次数都算用总订单数 ÷ 下单用户数。两个口径跑出来的数字能差 15%看板做完了才发现两边在吵架。所以正确顺序是先把指标的分子分母、时间窗口、去重维度写进一张表双方签字然后再动手写 Spark 脚本。具体要确认的东西包括时间字段用事件时间订单创建时间还是处理时间Spark 收到消息的时间窗口是滑动窗口还是滚动窗口窗口对不齐会不会导致数字跳变去重维度是 user_id 还是 device_id匿名用户怎么算迟到数据等多久超过阈值的直接丢还是补算这四条确认完Spark 侧的写法基本就定死了。反过来如果先写代码后对口径你大概率要重写一遍窗口逻辑。2. 环境落地Kafka 集群和 Spark 消费端的版本对齐2.1 Kafka 安装里最先绊人的几个细节Kafka 集群安装看起来简单解压改配置就行但有几个细节第一次做必踩。第一个是Windows 下启动脚本的路径问题。很多人复制命令的时候会把整个路径硬编码进去比如kafka-server-start.bat d:/rk/zy/kafka/kafka_2.13-3.0.0/config/server.properties。这样做的问题是脚本内部对相对路径的解析基准会变日志目录、数据目录都按当前工作目录去找结果就是启动成功但数据落到一个你找不到的地方。正确的做法是先cd到 Kafka 根目录再用相对路径cd /d d:\rk\zy\kafka\kafka_2.13-3.0.0 bin\windows\kafka-server-start.bat .\config\server.properties第二个是Docker 安装 Kafka 时的监听器配置。这是新手最大的坑没有之一。Kafka 有三个 listener 相关配置listeners、advertised.listeners、listener.security.protocol.map。listeners是 broker 实际绑定的地址advertised.listeners是 broker 告诉客户端你应该用这个地址来找我的地址。如果你在容器里把 advertised 配成localhost:9092那么宿主机之外的任何客户端拿到这个地址都会连不上——因为它会去连自己的 localhost。单机测试用 docker-compose 起 Kafka 时至少要保证 advertised 的地址是客户端能路由到的地址。跨主机访问时Kafka 不解决网络可达性问题它只负责告诉你去哪找我。第三个是伪集群。学习阶段不想开三台机器可以在一个节点上启动三个 broker用不同的broker.id、port、log.dirs。这里要注意log.dirs必须分开而且zookeeper.connect指向同一个 ZK。三个 broker 起来之后用kafka-topics.sh --create --replication-factor 3建 topic才能真正验证副本机制。副本数超过 broker 数会直接报错这也是新手常犯的。2.2 Topic 分区数和副本数怎么估分区数有两个作用决定生产端能并行写多少也决定消费端能并行读到多少。Structured Streaming 里一个 Kafka 分区对应一个 Spark task在初始读取阶段分区数少于你想要的并行度你的集群再大也只能干瞪眼。估算过程大致是这样假设你的峰值写入是 20 万条/秒单分区实测能稳定承载 3 万条/秒那么分区数至少是 7取整到 8。然后再看消费侧如果你给这个 streaming 作业分配了 8 个 executor core那 8 个分区刚好一一对应不多不少。有几个经验值可以参考分区数宁多勿少。Kafka 的分区只能增加不能减少增加分区会导致相同 key 的消息落到不同分区顺序性被破坏。所以一开始就要留出 2~3 倍的冗余。副本数生产环境给 3同时把min.insync.replicas设为 2。这样允许挂一台 broker 仍然可写挂两台才停止写入。消息 key 决定分区。如果你的下游要按 user_id 做有状态聚合就用 user_id 做 key如果只做全局统计不指定 key 让消息轮询分布反而更均匀。Topic 建好之后记得验证直接用命令行看./kafka-topics.sh --bootstrap-server localhost:9092 --describe --topic user-behavior输出里能看到每个分区的 leader、replicas、isr。ISR 列表长度小于副本数说明有副本没跟上这种情况下去查消费延迟是白费力气。2.3 Spark 和 Kafka connector 的版本对应关系这块是版本地狱的重灾区。Structured Streaming 读 Kafka 用的是spark-sql-kafka-0-10这个包包名里的 0-10 是老的 Kafka 客户端 API 版本但它对 Kafka 0.10 以上的版本都能用。真正的坑在于Spark 版本和 connector 版本必须严格对应。Spark 版本connector 坐标说明2.4.xspark-sql-kafka-0-10_2.11:2.4.8Scala 2.11 和 2.12 是两个包3.0 ~ 3.2spark-sql-kafka-0-10_2.12:3.1.33.0 起 Scala 2.12 成为主流3.3 ~ 3.5spark-sql-kafka-0-10_2.12:3.5.1功能稳定推荐注意坐标里的_2.11/_2.12是 Scala 版本后缀它必须和你 Spark 发行版的 Scala 版本一致不一致会在运行时报NoSuchMethodError而且报错信息完全看不出是版本问题。判断方法很简单看 Spark 安装包的文件名spark-3.5.1-bin-hadoop3后面如果没写 Scala 版本就看jars目录里scala-library的版本号。2.4 引入依赖的三种方式各自适合什么场景方式命令优点缺点--packages自动从中央仓库拉省事版本清晰需要外网内网离线环境不可用--jars手动下 jar 指定路径内网可用依赖传递要自己解决fat jar打成一个 uber jar依赖封闭不污染集群和集群自带的包容易冲突内网环境我一般用--jars加一份手工整理好的依赖清单。这里有个细节spark-sql-kafka-0-10依赖kafka-clients和commons-pool2如果你只下了主包运行时会报ClassNotFoundException: org.apache.kafka.clients.consumer.ConsumerConfig。解决办法是用一个干净的目录mvn dependency:copy-dependencies把整个依赖树拉下来一次性全丢进--jars。另外提醒一句不要把 connector 的 jar 直接放进 Spark 的 jars 目录。集群上如果有多个作业用不同版本的 connector这会导致互相覆盖排查起来非常痛苦。用--jars让每个作业自带依赖是更干净的做法。3. 消费端的坑offset、checkpoint 和重复消费的死结3.1 一次重启后数据翻倍的完整排查链路回到开头朋友那个问题。数据为什么会重复我把当时的排查顺序完整记下来这个顺序可以复用。第一步确认重复的范围。是某个时间段的数字整体翻倍还是只有某几个分区翻倍。当时查下来是整点那几个批次翻倍说明问题出在批处理边界不是单条消息。第二步看消费组的 offset。用命令行看消费组的当前位点./kafka-consumer-groups.sh --bootstrap-server localhost:9092 \ --describe --group dashboard-stream输出里的CURRENT-OFFSET是当前提交的位点LOG-END-OFFSET是分区最新位点两者之差就是 lag。当时看到 lag 正常说明消费组在正常工作那重复的数据就不来自重复消费。第三步看输出端的写入方式。把落库代码翻出来一看用的是mode(append)直接往聚合表里插。问题基本就浮出来了如果 Spark 作业重启后重放了上一批数据append 会把同一批数据再插一遍。第四步看 checkpoint 目录。这才是根因。他们的 checkpoint 目录配在一个本地临时盘上容器重启后这个盘被清了。checkpoint 丢了的后果是Spark 失去了我处理到哪了的记忆但它同时又通过startingOffsets或者消费组 offset 拿到了一个位点两边的位点对不上就出现了重放。更糟的是重放的那批数据写进聚合表时用的是 append于是翻倍。第五步验证。把 checkpoint 目录换成可靠存储把落库改成幂等写入重新上线跑了三天数字稳定。这个链条里最关键的一步是第三步和第四步。很多人卡在第二步看到 lag 正常就以为消费端没问题但 lag 只反映没消费到的量完全不反映消费了但重复写了。3.2 checkpoint 和消费组 offset到底谁说了算这个问题值得单独说清楚因为它决定了你所有重启行为的表现。Structured Streaming 的 Kafka source 有两种位点管理方式checkpoint 优先当 checkpoint 目录里有有效的 offset 记录时Spark 会直接从 checkpoint 里恢复位点startingOffsets参数完全不生效。startingOffsets 生效只在 checkpoint 目录为空首次启动或目录被清时才起作用。这意味着什么意味着只要 checkpoint 还在你把startingOffsets改成earliest也不会重跑历史数据。反过来说checkpoint 一丢无论你之前跑得多稳都得从头决定位点。所以生产环境的两条硬规矩checkpoint 必须放在可靠、持久、所有 executor 都能访问的存储上比如 HDFS 或者对象存储。绝不能放本地盘也绝不能用local[*]模式的默认路径。checkpoint 和消费组不要混用。如果你在 Spark 里指定了group.idSpark 确实会用消费组来提交 offset但它的恢复逻辑仍然优先看 checkpoint。两套机制并存只会让行为变得难以预测建议关掉enable.auto.commit让 checkpoint 做唯一的真相来源。还有一个容易忽略的点checkpoint 目录里保存的不只是位点还有查询计划的序列化信息。所以当你修改了作业的计算逻辑比如加了一个聚合字段重新提交时可能会报反序列化错误。这种情况下的标准做法是换一个新的 checkpoint 目录让它从头开始跑。代价是要手工补一下历史位点但改逻辑再来一次更贵。3.3 把 at-least-once 变成业务上的 exactly-onceKafka Spark 这套组合端到端默认是at-least-once语义可能会重复但不会丢。想做到真正的 exactly-once成本非常高需要事务性写入 幂等读取对 sink 有硬要求。实际项目中我更推荐的做法是承认至少一次然后在写入端做幂等。这样既不用改架构也能保证看板上的数字是对的。具体有三种常见方案方案一基于主键的 upsert。把聚合表的主键设计成(指标名, 时间窗口, 维度值)写入时用REPLACE INTO或者INSERT ... ON DUPLICATE KEY UPDATE。同一个批次重复写结果一样。这个方案的适用前提是聚合结果本身可覆盖不适合累加型指标。方案二先写临时表再合并。用foreachBatch每个批次先把结果写到一个带 batch_id 的临时表再用一条 SQL 把它合并进正式表合并完成后删掉临时记录。这个方案能处理累加型指标代价是每次多一次写入。方案三结果表带版本号。每次写入带上 batch_id查询时只取每个主键下最大的 batch_id。这个方案的优点是保留历史方便回溯缺点是存储会膨胀需要定期清理。def upsert_batch(batch_df, batch_id): # 先写临时表主键含 batch_id (batch_df.write .format(jdbc) .option(url, JDBC_URL) .option(dbtable, tmp_realtime_kpi) .mode(append) .save()) # 再按主键覆盖进正式表 conn get_conn() cur conn.cursor() cur.execute( REPLACE INTO ads_realtime_kpi SELECT metric, win_start, dim_value, metric_value FROM tmp_realtime_kpi WHERE batch_id %s , (batch_id,)) cur.execute(DELETE FROM tmp_realtime_kpi WHERE batch_id %s, (batch_id,)) conn.commit()这段代码里有几个细节临时表和正式表要在同一个库里否则 REPLACE 的跨库事务会有性能问题batch_id要建索引删除时才不会全表扫整个操作要放在一个事务里避免合并到一半进程挂掉留下脏数据。3.4 lag 高的时候先看这四个数Kafka 消息延迟高很多人第一反应是加机器。先别急看四个数。观察项命令/位置正常表现异常时的含义分区 lag 分布consumer-groups --describe各分区 lag 接近单个分区高说明数据倾斜ISR 长度topics --describe等于副本数小于副本数说明有副本同步慢生产端速率broker 的 MessagesInPerSec平稳突然飙升说明上游有活动消费端速率消费组 offset 变化量跟得上生产跟不上就是消费能力不足判断逻辑是这样的如果 lag 在涨但各分区均匀那是整体消费能力不足方向是加并行度分区数够的话加 executor core不够就得先扩分区。如果只有个别分区 lag 高那是数据倾斜要去看生产端是不是用了不均匀的 key比如某个大客户的 user_id 占了 40% 的消息量。还有一种更隐蔽的情况lag 不涨但看板数据在延迟。这通常不是 Kafka 的问题而是 Spark 的 trigger 间隔 落库延迟 Dashboard 轮询周期三者叠加出来的。10 秒 trigger 2 秒落库 30 秒轮询最坏情况下数据要 42 秒才反映到看板上。要缩短这个链路得从每一段去抠。4. Spark 侧的计算和内存调优4.1 几个必须显式设置的参数Structured Streaming 读 Kafka 的代码就那么几行但默认值大多数不适合生产。raw (spark.readStream .format(kafka) .option(kafka.bootstrap.servers, kafka1:9092,kafka2:9092,kafka3:9092) .option(subscribe, user-behavior) .option(startingOffsets, latest) .option(maxOffsetsPerTrigger, 200000) .option(failOnDataLoss, false) .option(kafka.isolation.level, read_committed) .load())逐个说下为什么。maxOffsetsPerTrigger是限流阀。默认不设的话Spark 会把当前所有可用数据一次性拉过来如果作业停了半小时再启动第一个批次可能就是几千万条直接 OOM。设成 20 万相当于给每个批次上了天花板处理不过来的部分留在 Kafka 里下一批继续。这个值怎么定用你的单批次处理耗时反推如果 20 万条能在 5 秒内处理完trigger 设 10 秒那这个值就是合理的。failOnDataLoss默认为true意思是如果 Kafka 里的数据因为过期被删了、导致 Spark 想读的位点不存在了作业直接失败。生产环境建议设成false让作业跳过缺失的数据继续跑。但要配合监控——如果这个选项触发得太频繁说明你的消费速度已经跟不上日志保留速度了得调整retention.ms或者提高消费能力。isolation.level设成read_committed是为了配合上游的事务性生产者。如果上游用的是普通生产者这个选项没影响如果上游开了事务不设这个会读到未提交的消息然后在下游出现数据凭空消失的怪现象。4.2 内存模型和 OOM 的真实成因Spark 的内存分三块executor 内存堆内、executor 内存开销堆外默认是堆内的 10%、driver 内存。OOM 报错里的关键词决定了你要调哪一块。java.lang.OutOfMemoryError: Java heap space出现在 executor 日志里 → 调spark.executor.memory同样报错出现在 driver 日志里 → 调spark.driver.memoryContainer killed by YARN for exceeding memory limits→ 是堆外超了调spark.executor.memoryOverhead有一个很常见的误区executor 内存不是越大越好。单个 executor 给到 32G 以上GC 停顿会明显变长反而拉低吞吐。我一般控制在 8G~16G 之间然后用增加 executor 数量的方式扩容。还有一个参数值得单独提spark.memory.fraction默认 0.6表示堆内存中用于执行和缓存的比例。做有状态聚合比如窗口去重的作业状态会占用大量内存如果发现 GC 频繁但堆没满可以适当调到 0.7。反过来如果作业里有大量自定义对象把它降到 0.5 给用户对象留空间。另外状态存储默认在内存里。窗口去重如果窗口开得大比如 7 天去重状态会持续膨胀。这时候要开 RocksDB 状态存储spark.sql.streaming.stateStore.providerClassorg.apache.spark.sql.execution.streaming.state.RocksDBStateStoreProviderHDFS 上的 checkpoint 会承载这部分状态速度比纯内存慢但不会 OOM。4.3 分区数、shuffle 分区、限流量三者的关系这三个数经常被混在一起谈其实它们管的是不同阶段参数作用阶段影响的并行度Kafka 分区数读取初始 batch 的 task 数spark.sql.shuffle.partitions聚合/joinshuffle 后的 task 数maxOffsetsPerTrigger读取单批次数据量上限spark.sql.shuffle.partitions默认 200这个默认值在两种场景下都不合适数据量小的时候200 个 task 每个只处理几百条调度开销比计算还大数据量大的时候200 个 task 每个处理上百万条又会 OOM。判断方法很简单看 Spark UI 里每个 task 的处理时间。如果有一批 task 几毫秒就结束、另一批要几十秒说明分区不均如果所有 task 都要几十秒且内存高说明分区太少。目标状态是每个 task 处理 128MB 左右的数据耗时在几秒量级。对于流式作业shuffle.partitions没法像批处理那样自适应。我的做法是先跑一遍历史数据全量回填从 Spark UI 拿到理想的分区数然后把这个值固定到作业配置里。数据量波动大的话可以在代码里根据 batch 的行数动态设置batch_df ... num_partitions max(4, min(200, int(batch_df.count() / 100000))) spark.conf.set(spark.sql.shuffle.partitions, str(num_partitions))4.4 时间窗口和 watermark迟到数据该不该等窗口聚合必配 watermark否则状态永远不释放跑几天就 OOM。watermark 的含义是系统愿意等待迟到数据的最大时长。设成 10 分钟意思是如果一个 10:00 的数据在 10:10 之后才到它会被丢弃10:00 这个窗口在 10:10 之后不再更新。agg (raw .withWatermark(event_time, 10 minutes) .groupBy(window(event_time, 5 minutes), dim_value) .agg(count(*).alias(cnt)))这里有两个非常容易混淆的点。第一窗口能不能对上。5 分钟窗口的起点是固定的10:00-10:05、10:05-10:10不会因为 watermark 而变化。这一点对看板很重要——如果窗口起点动态变化仪表盘上的历史点会不停被改写用户会觉得数字在跳。第二watermark 设多大。太小迟到数据被丢看板数字偏低太大状态保留时间长内存压力大。我的经验值是取迟到分布的 P99 再加一点余量。怎么拿到迟到分布在数据里加两个字段事件时间和处理时间跑一天之后统计两者差值的分布P99 是 3 分钟的话watermark 设 5 分钟就够了。消费端的真实场景里迟到往往不是网络问题而是上游的业务延迟。比如移动端在网络差的时候会缓存埋点等有网了再批量上报。这种上报的延迟可能到几小时甚至跨天。如果你确实需要这部分数据watermark 就得开到很大或者干脆走两条链路实时链路只处理准时数据全量链路第二天补算。5. Dashboard 展示层让数字不跳动、不打架5.1 预聚合表怎么设计看板永远查预聚合表绝不查明细。这一点没有讨论余地。聚合表的粒度设计遵循最小可用维度集合原则。假设看板上有这几个需求按小时看销售额趋势、按渠道看销售额占比、按商品类目看 Top10。那么最小维度集合是(小时, 渠道, 类目)每一行是这三个维度组合下的销售额。任何时候的查询都是在这张表上按需过滤和再聚合。表结构大致是这样CREATE TABLE ads_realtime_kpi ( metric_code VARCHAR(64) NOT NULL COMMENT 指标编码, win_start DATETIME NOT NULL COMMENT 窗口起点, dim_channel VARCHAR(32) NOT NULL DEFAULT all, dim_category VARCHAR(32) NOT NULL DEFAULT all, metric_value DECIMAL(20,4) NOT NULL DEFAULT 0, updated_at TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP, PRIMARY KEY (metric_code, win_start, dim_channel, dim_category) ) ENGINEInnoDB DEFAULT CHARSETutf8mb4;几个设计要点dim_*字段用all而不是 NULL 表示不区分维度。因为 MySQL 的唯一索引里 NULL 不参与去重用 NULL 会让 REPLACE INTO 失效这是幂等写入最常见的失败原因。win_start用窗口起点而不是窗口结束时间。起点在窗口未闭合时就已经确定了可以提前写入用结束时间的话未闭合的窗口没法写看板就会看到数据延迟一个窗口。metric_value用 DECIMAL 不用 FLOAT。金额计算用浮点数累积误差会让看板的总和对不上明细运营一定会来问。加updated_at便于排查。看板上数字没更新时先看这个字段的时间戳能立刻区分是 Spark 没写还是 Dashboard 没读。5.2 刷新节奏和窗口对齐看板数字跳动除了数据重复还有一个常见原因是刷新节奏和窗口边界没对齐。举个具体的例子窗口是 5 分钟Spark 的 trigger 是 10 秒Dashboard 每 30 秒轮询一次。那么在一个 5 分钟窗口内这个窗口的数值会被更新大约 10 次——每次都是到目前为止的部分和。用户看到的是数字每 30 秒往上涨一点一直涨到窗口闭合。这本身没问题但如果有两个窗口同时在屏幕上当前窗口和上一个窗口就会出现前面那个不涨、后面那个在涨的视觉效果用户很容易以为是数据错了。解决办法有两个方案 A只展示已闭合的窗口。查询条件加win_start 当前时间 - 窗口长度牺牲一个窗口的延迟换取数字稳定。这是最省事的做法适合对延迟不敏感的业务看板。方案 B区分进行中和已闭合。用两种视觉样式呈现进行中的用虚线或浅色闭合的后用实线。这样用户既能感知实时性也知道哪个数字还会变。这个方案体验更好但前端要多做点事。还有一种情况值得提如果用到累计值类指标比如今日累计 GMV要用窗口求和而不是直接读某个窗口。因为窗口会重算直接读单窗口的值在窗口闭合瞬间会有一次跳变。5.3 可视化工具怎么选这块没有标准答案看团队的技术栈和需求复杂度。我的判断逻辑是工具适合的场景不适合的场景自研前端 ECharts交互复杂、要嵌入自有系统团队没有前端资源Superset自助分析、维度组合多强实时、秒级刷新Grafana监控类、时序类指标复杂维度下钻大屏工具演示、汇报、固定布局日常运营使用需要提醒的是不要用监控工具的思路做业务看板。Grafana 擅长的是一个指标在不同时间点的值业务看板往往需要多个指标在同一时间点的对比 维度下钻。强行用监控工具做业务看板最后会变成几十个 panel 堆在一起没人看得懂。还有一个实操建议看板上的每个数字都要能点进去看到明细。哪怕只是跳转到一个预聚合更细粒度的页面。运营看到数字的第一反应永远是这个对不对如果不给他追溯的路径他会直接来找你你的时间就被无限切碎。6. 踩坑笔记复盘三张排查清单6.1 数据不更新了按这个顺序查第一层链路是否活着。先确认三件事Kafka 有没有新消息进来用命令行看 topic 的 offset 有没有涨、Spark 作业是不是还在跑看 Spark UI 的 Streaming 页面、结果表的updated_at有没有变化。这三步能在两分钟内定位是哪一段断了。第二层读到数据没有。如果 Kafka 有数据、Spark 也在跑但结果表不更新看 Spark UI 里每个 batch 的numInputRows。如果一直是 0说明读不到数据——检查subscribe的 topic 名拼写、检查消费组的位点是不是停在了一个已经没有数据的位置。第三层是不是被过滤掉了。numInputRows有值但输出没有通常是代价条件的问题。常见的是 watermark 设太小导致所有数据都被判为迟到或者时间字段解析失败变成 null被过滤掉了。这里有个小技巧在流式 DataFrame 上加一个foreachBatch打印批次行数能快速区分没读到和读到但没输出。第四层写入是不是静默失败了。JDBC 写入失败时如果异常被吞掉作业会继续跑但数据不落库。所以foreachBatch里的异常一定要往外抛让整个批次失败并重试。# 快速确认 Kafka 里到底有没有数据 ./kafka-console-consumer.sh --bootstrap-server localhost:9092 \ --topic user-behavior --from-beginning --max-messages 56.2 数据重复了按这个顺序查先确认重复的形态。整体翻倍、局部翻倍、还是有零散的重复行对应的根因完全不同。整体翻倍通常有一个明确的时间点就是作业重启。查 checkpoint 目录是否存在、是否被清空、重启时有没有切换 checkpoint 路径。局部翻倍通常和某个分区有关。查上游是不是有重发逻辑生产端的重试 没有开幂等 消息重复这类重复是 Kafka 层面的需要生产端开enable.idempotencetrue解决。零散重复最常见于 join 场景。流式 join 静态表如果静态表里有重复的 keyjoin 后行数会膨胀。这种情况要在 join 之前对静态表按 key 去重或者用dropDuplicates。排查的时候有个通用手法给每一行打上 batch_id 和输入 offset重复数据能直接从这两个字段看出是同一批次重放还是不同批次重复。这个字段只在排查期间加稳定之后可以去掉避免额外的存储开销。6.3 那些改了就有明显效果的默认值这些参数我整理成了一张表搭新作业时直接抄参数默认值建议值原因spark.sql.shuffle.partitions200按数据量算通常 20~100默认值两头不讨好spark.sql.streaming.stateStore.providerClass内存RocksDB大状态场景必开spark.serializerJavaKryo序列化开销降低明显maxOffsetsPerTrigger无按处理能力估算防止首批发车过载failOnDataLosstruefalse避免因数据过期而整作业挂掉spark.executor.memoryOverhead堆的 10%堆的 15%~20%容器被杀多半是这里超了trigger interval无尽快10s / 30s减少小批次开销稳定输出节奏Kryo 序列化有个坑要提醒它要求注册自定义类才能发挥全部性能不注册也能用但会退化成把类名写进流里效果打折。如果作业里没有自定义类用 Kryo 就够了如果有记得在sparkConf里注册。最后说一个跟技术无关但很重要的经验。这套链路搭完之后我建议先在测试环境用真实流量的回放跑满一周再上生产。回放的方式很简单把生产流量录一份到测试 topic用相同的 Spark 作业消费。一周时间里你会遇到上游业务的各种活动峰值、会遇到 Kafka 的日志滚动、会遇到自动扩缩容这些问题如果在上生产当天暴露代价会高很多倍。另外看板上线的第一天一定要有人盯着。不是盯技术指标是盯业务方看到数字之后的反应。他们会问为什么这个数跟我算的不一样这个时候能立刻定位是口径问题还是技术问题比事后翻日志高效得多。这套流程我走过几次系统真正稳下来往往不是因为技术改得多好而是因为口径终于对齐了。
返回列表