Flink 实时开发面试深度准备(AviaGames · 数据工程师)
本文分两部分:第一部分是知识点细节(按"面试官会怎么追问"的链条展开,以布隆过滤器开头),第二部分是项目细节(把你手上的真实排错经历改写成面试版 STAR 叙述,并预演追问)。
第一部分:知识点细节(追问链模式)
1. 布隆过滤器(上次没答上来的,这次答透)
1.1 一句话原理
一个长度为 m 的 bit 数组 + k 个独立哈希函数。插入时把 k 个位置置 1;查询时 k 个位置全为 1 才认为"可能存在",任何一位为 0 则"一定不存在"。
1.2 面试官问"布隆过滤器有哪些问题"——标准答案(按重要性排序)
① 假阳性(False Positive),且不可消除
- 只会"误判存在",不会"误判不存在"。原因:不同 key 的哈希位会重叠,查询的 key 命中的可能全是别人置的位。
- 误判率公式:p ≈ (1 - e^(-kn/m))^k,其中 n 是已插入元素数。
- 在去重场景的致命后果:误判"存在" = 把一条新数据当成重复数据丢掉 = 丢数据。所以 BF 去重只适合允许少量丢失的场景(爬虫 URL 去重、缓存穿透防护、推荐去重),绝对不能用于金融交易、对局结算这类要求精确的链路。这一句是资深和初级的分水岭,一定要主动说。
② 不支持删除 - 一个 bit 可能被多个 key 共享,删除时置 0 会"误伤"其他 key,导致假阴性(这比假阳性更不可接受)。 - 解决方案:Counting Bloom Filter,把每个 bit 换成计数器(通常 4 bit),删除时减 1。代价:空间膨胀 4 倍,且计数器有溢出风险。
③ 容量必须预估,不能动态扩容 - m 和 k 在创建时就定死。实际元素数超过预估 n 后,误判率会急剧恶化(位数组被填满,趋近于全 1)。 - 解决方案:Scalable Bloom Filter——满了就再挂一个新的、误判率更严格的 BF,形成链。代价:查询要查所有层,越用越慢。 - 实践中通常直接预留 1.5~2 倍容量,或按天/按窗口滚动重建(实时去重场景最常用:每天 0 点换一个新 BF)。
④ 无法取出原始数据、无法遍历、无法统计精确数量 - 它只回答"in or not in",不能列出里面有什么。需要精确对账时没法用。
⑤ 哈希计算成本与哈希质量 - 每次操作要算 k 次哈希。工程上用 double hashing 技巧优化:只算两个哈希 h1、h2,第 i 个哈希 = h1 + i×h2,效果接近 k 个独立哈希。 - 哈希分布不均会让实际误判率高于理论值。
⑥ 分布式合并的约束 - 两个 BF 可以按位 OR 合并(用于多分区/多 subtask 局部去重后合并),但前提是 m、k、哈希函数完全一致,否则不能合并。
1.3 追问:参数怎么定?(会考计算)
给定预估元素数 n 和目标误判率 p:
每元素位数:m/n = -ln(p) / (ln2)²
最优哈希个数:k = (m/n) × ln2
背一个实例:1 亿个 key、1% 误判率 → m/n ≈ 9.6 bit/key → 总共约 114 MB,k ≈ 7 个哈希函数。对比:1 亿个 64 字节的 key 直接存 HashSet 至少 6.4 GB+,这就是 BF 的价值——用 ~2% 的空间换 1% 的误判。
误判率每降一个数量级(1% → 0.1%),大约多花 4.8 bit/key,空间是线性温和增长的,所以"要不要更低误判率"是个很划算的权衡题。
1.4 追问:有什么替代品?什么时候不用 BF?
| 方案 | 精确性 | 支持删除 | 空间 | 适用 |
|---|---|---|---|---|
| Bloom Filter | 有假阳性 | 否 | 极省 | 允许丢的去重、存在性预判 |
| Counting BF | 有假阳性 | 是 | BF×4 | 需删除的近似去重 |
| Cuckoo Filter | 有假阳性 | 是 | 低误判率(<3%)时比 BF 更省 | BF 的现代替代,查询只碰 2 个桶,cache 友好;缺点是接近满载(~95%)时插入可能失败/踢出循环 |
| HyperLogLog | 误差 ~0.8% | — | 12 KB 固定 | 只做基数统计(UV 数),不能判断单条是否重复——很多人混淆,面试官爱挖这个坑 |
| Roaring Bitmap | 精确 | 是 | 取决于分布 | key 可整数化(用户 ID 字典编码)时的精确去重/交并集,Doris 的 bitmap 类型就是它 |
| Flink ValueState + TTL | 精确 | 是 | 大(每 key 一条 state) | 金融级精确去重,配 RocksDB 落盘 |
回答模板:"要不要用 BF,先问业务能不能容忍丢数据。能容忍→BF/Cuckoo;不能容忍但 key 能整数化→Roaring Bitmap;不能容忍且 key 任意→state 精确去重(RocksDB 扛量)+ TTL 控制状态大小;只要 UV 数字→HLL。"
1.5 追问:BF 在你熟悉的系统里哪里出现?
- RocksDB(也就是 Flink 的 state backend):每个 SST 文件内置 BF,点查时先过 BF 决定要不要读这个文件,大幅减少读放大。
state.backend.rocksdb调优里 bloom filter bits 是一项。 - HBase / Doris:BF 索引用于加速点查的"负缓存"——快速判定这个 block/tablet 里没有目标 key,跳过读盘。Doris 建表可指定
bloom_filter_columns。 - Redis 缓存穿透防护:请求先过 BF,不存在的 key 直接挡回,不打到 DB。
- Flink 实时去重:大流量下先用 BF 挡掉 99% 的"肯定没见过",剩下的再查精确 state——两级去重,兼顾内存和精确性(注意方向:BF 说"不存在"是可信的,说"存在"再去查精确层确认)。
2. Flink 去重方案全景(BF 的自然延伸,必考)
场景:对局事件可能重复上报(客户端重试),按 event_id 去重。
- ValueState<Boolean> + StateTtlConfig:精确,每个 key 一条 state。10 亿级 key 用 RocksDB backend 落盘。TTL 按业务窗口设(比如对局事件 24h 后不可能再重复 → TTL 24h),防止 state 无限膨胀。TTL 清理方式:
cleanupInRocksdbCompactFilter(compaction 时顺带清理,推荐)。 - Flink SQL:
ROW_NUMBER() OVER (PARTITION BY event_id ORDER BY proctime) WHERE rn = 1,配table.exec.state.ttl。注意 deduplication 保留第一条(append 流)和保留最后一条(retract 流)对下游的影响不同。 - BF/两级去重:上面 1.5 说过,大流量省内存,有丢数风险。
- 交给下游幂等:Flink 端 at-least-once,Doris 用 Unique Key 模型按主键覆盖——工程上最常用的折中,把去重成本转移给存储层。要会说这句:"端到端一致性不一定要在计算层硬扛,sink 幂等 + 上游至少一次,往往比 2PC 更简单可靠。"
3. Checkpoint 深水区(必考,追问能到很深)
3.1 原理链条(准备被连环追问)
- Chandy-Lamport 异步屏障快照:JobManager 触发 → Source 注入 barrier → barrier 随数据流动 → 每个算子收到 barrier 时对自己的 state 做快照 → 全部完成后 checkpoint 成功。
- Barrier 对齐(exactly-once 模式):多输入算子要等所有 channel 的 barrier 到齐,先到的 channel 数据会被缓存(1.11 前是阻塞 channel)。反压时 barrier 流得慢 → 对齐等很久 → checkpoint 超时。这是"反压导致 checkpoint 失败"的机制解释,必须能讲出来。
- At-least-once 模式:不对齐,先到先处理,快照包含了"越过 barrier"的数据,恢复时重放会重复。
- Unaligned Checkpoint(1.11+):barrier 直接越过 in-flight 数据插队,把通道中的 in-flight buffer 也写进 checkpoint。反压严重时 checkpoint 依然能快速完成。代价:state 体积变大、恢复变慢、不能配合 rescale 直接改并行度恢复(1.13 后放宽)。1.13+ 支持
alignment timeout:先尝试对齐,超时自动切换成非对齐——生产推荐这个混合模式。 - 增量 Checkpoint(RocksDB):只上传新生成的 SST 文件。代价:checkpoint 之间形成文件依赖链,单次恢复要下载的总量可能不小;长期运行需关注 compaction 是否跟得上。
- Changelog State Backend(1.15+):持续把 state 变更写 changelog,checkpoint 时只需刷很小的增量,把 checkpoint 时长从"分钟级"压到"秒级"。可以提一句表示跟进社区。
3.2 排查话术:"Checkpoint 老是超时/失败,你怎么查?"
标准五步,背熟:
- 看 Web UI Checkpoint 详情页,定位是哪个算子、哪个 subtask 慢——是所有 subtask 都慢(全局问题:存储慢、状态太大)还是个别慢(倾斜)。
- 拆分阶段:
Sync Duration(本地快照,通常快)、Async Duration(上传远端,慢=存储带宽/状态大)、Alignment Duration/Start Delay(大=barrier 传播慢=反压)。Start Delay 高几乎可以直接判反压。 - 反压路线:去查反压(见第 4 节),根因通常在下游算子或 sink。
- 状态大路线:看 state size 趋势是否无界增长 → 检查有没有裸的 regular join / 无 TTL 的 keyed state / 迟到数据导致窗口不释放。
- 兜底手段:开 unaligned/alignment-timeout、增大 checkpoint 超时与间隔(
min-pause防止背靠背)、增量 checkpoint、必要时 state 瘦身(TTL、改 interval join)。
3.3 高频追问点
- Checkpoint vs Savepoint:前者自动、面向容错、格式可增量;后者手动、面向版本升级/迁移、规范格式。任务升级要给算子固定 uid,否则 state 对不上。
- 恢复语义:重启后从最近 checkpoint 恢复,source 回拨 offset(offset 存在 state 里,不依赖 Kafka 提交的 offset——
enable.auto.commit和一致性无关,常见坑题)。 - Task Local Recovery:本地留一份快照副本,failover 时不用从远端拉,加速恢复。
4. 反压:机制 + 定位(他们的痛点区,重点)
4.1 机制(能讲出这段就赢一半)
Flink 1.5+ 是 credit-based 流控:下游每个 input channel 向上游通告自己的可用 buffer 数(credit),上游只在有 credit 时才发数据,并附带自己的积压量(backlog)供下游申请更多 buffer。所以反压是逐级、精确到 channel 的,不会像 TCP 反压那样把整个连接堵死。反压的表现 = 上游 output buffer 耗尽 → 上游也停 → 一路传导到 source(source 消费 Kafka 变慢,consumer lag 上涨——所以 lag 报警经常是反压的第一信号)。
4.2 定位五步法
- Web UI 作业拓扑着色 + 三个指标:
backPressuredTimeMsPerSecond/busyTimeMsPerSecond/idleTimeMsPerSecond。从 source 往下找第一个"busy 高、backpressured 低"的算子——它就是瓶颈;它上游全部会显示 backpressured。 - 确认瓶颈类型:开火焰图(
rest.flamegraph.enabled: true)或 async-profiler 看 CPU 热点——是业务逻辑重、序列化(Kryo fallback 是常见元凶,POJO/Tuple 快得多)、还是外部 IO(同步查维表/写 sink)。 - 看是不是倾斜:同一算子各 subtask 的
numRecordsIn/ busy 分布,个别 subtask 100% busy 其他闲 = 倾斜,走倾斜处理(第 5 节),加并行度没用。 - 看 GC:TaskManager GC 日志/metrics,老年代频繁 Full GC → 状态对象太多/内存配比不对(调 managed memory 比例,RocksDB 场景 managed memory 就是 RocksDB 的 block cache + write buffer)。
- sink 慢的专项:Doris stream load 批次太小太频繁(见第 7 节)、外部服务限流 → 加大 batch、异步化、扩 sink 并行度。
4.3 常见根因排名(可直接说"我按概率排查")
sink/外部系统慢 > 数据倾斜 > 序列化开销 > 单并行度瓶颈算子 > GC > 资源真的不够。
5. 数据倾斜(和反压绑定出现)
- 定位:subtask 级
numRecordsIn分布、各 subtask checkpoint state size 差异、Kafka 各分区流量(倾斜可能从 Kafka 分区 key 就开始了,比如按游戏房间 ID 分区而头部房间超热)。 - DataStream 处理:两阶段聚合加盐——
keyBy(key + rand(N))局部聚合 →keyBy(key)全局合并。热点 key 广播/单独拆流处理。 - SQL 处理(背参数):
- MiniBatch 三件套:
table.exec.mini-batch.enabled=true、allow-latency=1s、size=5000——攒批减少 state 访问次数; table.optimizer.agg-phase-strategy=TWO_PHASE:local-global 聚合,先本地预聚合再全局,等价于 MR 的 combiner;table.optimizer.distinct-agg.split.enabled=true:专治COUNT(DISTINCT user_id)热点,自动按 distinct key 打散成两层。- join 倾斜:小表广播(broadcast state / broadcast join);大表热 key 加盐 + 另一侧膨胀复制。
- 说一个真实教训(见第二部分故事 7):GROUP BY 把所有 NULL 当同一组,NULL key 群会伪装成"超大热点/超高重复度"——先
IS NOT NULL再谈倾斜。
6. 时间语义、Watermark、迟到数据(游戏事件流必考)
- Watermark 本质:一条特殊记录,声明"event time ≤ t 的数据(基本)到齐了",用来触发窗口。多输入取最小值。
- KafkaSource 是 per-partition watermark:每个分区独立生成再取 min。某分区没数据 → watermark 不推进 → 全链路窗口不触发,解法
withIdleness(Duration)把空闲分区标记为 idle 不参与 min。这是生产最高频事故之一,主动讲。 - Watermark Alignment(1.15+):多源/多分区消费速度差异大时,快的源会导致慢源数据在 join/window state 里堆积;对齐机制让快源暂停等慢源。追赶历史数据(backfill)时特别有用。
- 迟到数据三层防御:① watermark 减去容忍时间(
forBoundedOutOfOrderness(5s),覆盖网络乱序);②allowedLateness(1h):窗口触发后不销毁,迟到数据来了增量再触发——注意下游必须能接受更新(retract/upsert sink),且窗口状态多保留 1h;③sideOutputLateData:兜底旁路,落到补偿表,离线修正。 - 游戏场景专属观点(主动说,很加分):客户端时间不可信——玩家改本地时间是常见作弊手段,event time 必须以服务端接收时间或服务端生成的对局时间为准。我在现在公司处理 MT 交易服务器时间时就立了同样的规矩:MT 服务器时间(GMT+2/+3 夏令时)是唯一权威时间,所有下游不得自行换算——时间口径不统一是数据不一致的头号来源。
7. 端到端 Exactly-Once 与 Flink→Doris(结合他们的栈,必考)
7.1 通用框架
端到端 EO = 可重放的 source(Kafka offset 在 state 里)+ 算子内 EO(checkpoint)+ 事务性或幂等 sink。三选一都不能少。
7.2 Kafka 事务 sink 的坑(资深题)
- 2PC 流程:checkpoint 时 preCommit(flush 到事务),
notifyCheckpointComplete时 commit。 - 坑 1:
transaction.timeout.ms必须 > checkpoint 间隔 + 最大 checkpoint 时长,且不能超过 broker 的transaction.max.timeout.ms(默认 15min)——否则任务恢复时事务已被 broker 中止,数据丢失。 - 坑 2:下游用
read_committed消费时,未提交事务会卡住 LSO,checkpoint 间隔就是下游可见延迟的下限——EO 是拿延迟换的。 - 坑 3:notifyCheckpointComplete 不是原子的,极端情况(commit 前挂掉)靠恢复后重放 pending transaction 补提交。
7.3 Flink → Doris(他们大概率在用,重点背)
- 方案 A:Connector 2PC(
sink.enable-2pc=true):stream load 分 precommit/commit 两段,label 与 checkpoint 绑定(jobId + checkpointId + subtask),label 全局唯一 = 天然幂等——重启重试同一批次,Doris 见到重复 label 直接拒绝,不会重复导入。 - 方案 B(工程上更常用):Flink at-least-once + Doris Unique Key 模型按主键覆盖。1.2+ 的 Merge-on-Write 让 Unique 表读时无需归并,点查/聚合性能接近 Duplicate 表。回答时给权衡:"结算金额类我用 2PC 或主键幂等兜底;PV 类指标 at-least-once 就够,不为它付 2PC 的延迟税。"
- Doris 侧高频事故(说出来就是加分项):stream load 批次太小太频繁 → 版本堆积 compaction 跟不上 → 报
-235 too many tablet versions。缓解:加大 Flink sink 攒批(sink.buffer-flush.max-rows/interval)、Doris 2.1+ 开 group commit、控制分桶数。 - Doris 三模型一句话:Duplicate(明细,啥都不干)、Aggregate(预聚合,SUM/MAX/REPLACE,排行榜适用)、Unique(主键去重,幂等 sink 的搭档)。
8. 双流 Join、维表关联(必考第二梯队)
- Regular join 的原罪:两侧全量进 state 且默认永不清理 → state 无限增长 → checkpoint 越来越大直到任务崩。必须设
table.exec.state.ttl,但要讲清语义风险:TTL 过期后来了本该 join 上的数据会 join 不上/输出错误结果——TTL 是正确性和资源的交易,要按业务最大关联跨度定(比如"平仓最晚发生在开仓后 30 天"→ TTL 31 天。我在现司做过统计:99.97% 的持仓在 30 天内平仓,以此把回看窗口从 90 天压到 30 天——同一个方法论)。 - Interval join:
a.time BETWEEN b.time - 10min AND b.time + 10min,state 随 watermark 自动清理,能用就优先用。 - 维表四方案:① broadcast state(小维表全量广播,更新走广播流);② Lookup join + cache(简单,cache TTL 内有脏读);③ Async IO + Redis/HBase(高吞吐,注意有序性
orderedWaitvsunorderedWait);④ 预加载 + 定时刷新(维表小且变更不频繁)。每个都要能说一致性/延迟/吞吐的权衡。 - 对局超时检测(他们业务的典型题):
KeyedProcessFunction按 match_id keyBy,收到"开局"事件存 state 并registerEventTimeTimer(开局时间 + 最大对局时长);收到"结算"事件删 timer 清 state;timer 触发 = 超时未结算,输出异常记录。追问点:timer 也在 checkpoint 里(不丢);海量 timer 用 RocksDB 存储(默认);processing time timer 和 event time timer 的触发条件差异。
9. Kafka 细节(快问快答级)
- 分区内有序、分区间无序 → 对局事件按 match_id 做分区 key 保证单场比赛事件有序;代价是热门时段头部倾斜,需在 key 设计上权衡。
- 可靠性三件套:
acks=all+min.insync.replicas=2+ 副本数 3;unclean leader election 关闭。 - Rebalance:新版 cooperative-sticky 协议增量 rebalance,避免 stop-the-world;Flink KafkaSource 自己管理分区分配和 offset(存 state),不依赖 consumer group rebalance——这也是为什么 Flink 任务的 offset 不能只看
kafka-consumer-groups命令。 - 分区数 vs Flink 并行度:并行度 > 分区数时多出的 subtask 空跑且(不设 idleness 时)拖住 watermark。
第二部分:项目细节(STAR 面试版 + 追问预演)
每个故事按"60 秒主叙述 + 追问预演"组织。主叙述要练到脱口而出,追问答案理解即可。数字全部要说——量级就是可信度。
故事 1:湖仓写入任务 hang 死的排查(展示:从现象到 root cause 的完整定位链)
主叙述: "我负责把 MT5 交易数据从旧架构迁到 Paimon 湖仓,单表 8 亿多行。首次全量写入时任务 hang 死:所有 executor 已空闲但作业不结束,不报错、不退出。我先排除了资源问题——executor 都空转,说明卡的不是计算阶段;然后看 driver 侧,发现卡在 Paimon 的 commit 阶段,commit 元数据大得异常。顺着查发现分区数量是百万级,而这张表按业务日期分区最多几千个分区。Root cause 是 DDL 里我把分区列写在了最前面,导致写入时列序错位,一个 md5 主键列被当成了分区键——每个 md5 值生成一个分区,几百万分区把 commit 的 manifest 元数据撑爆了。把分区列移到 DDL 末尾后,829M 行一次写入成功。之后我把'分区列必须放 DDL 末尾'写进了团队 DDL 规范。"
追问预演: - "为什么分区多 commit 会爆?" → 湖格式每次提交要生成 snapshot + manifest 记录本次涉及的所有分区/文件变更,百万分区意味着 manifest 条目爆炸,元数据序列化和提交本身成为瓶颈。同理 Hive 小文件多、Iceberg manifest 膨胀都是一类问题——元数据规模是湖仓的隐形瓶颈。 - "和 Flink 有什么关系?" → Paimon 就是 Flink 社区孵化的流式湖仓,我因此深入研究过它的 snapshot/changelog/tag 机制,对 Flink 流式写湖、CDC 语义(changelog-producer、first-row/lookup/full-compaction 的区别)有一手经验。 - "怎么避免再发生?" → DDL 规范 + 上线前小样本试写 + 监控单次 commit 的分区触达数。
故事 2:线上数据 NULL 回归,用快照 diff 定位到任务间竞态(展示:JD 第 5 条"迅速定位线上数据问题")
主叙述:
"小时级的平仓宽表突然出现 open_time 为 NULL 的记录。这张表的 open_time 来自和持仓表的关联。我用湖仓的 time-travel 能力(VERSION AS OF snapshot_id)对 NULL 出现前后的快照做 diff,先把问题记录锁定到具体的写入批次和时间点;再对照两个任务的调度时间线,发现是每日全量覆盖的持仓表任务和每小时增量的平仓表任务之间存在竞态——平仓任务恰好在持仓表'已清空、未写完'的窗口内跑了关联,自然关联不上。更深一层的结构性原因是上游数据保留窗口的设计问题。短期修复是调度依赖对齐,长期是消除覆盖写的中间不一致态。"
追问预演:
- "实时链路里怎么防这类问题?" → 这本质是维表新鲜度/一致性问题,Flink 里对应维表关联的 cache 一致性:用 broadcast 流推变更而不是定时全量替换、或 lookup cache 设短 TTL、或 temporal join 用版本表按事件时间关联到"当时"的维度值——temporal join 是根治手段。
- "time-travel diff 具体怎么做?" → 两个快照各取问题主键集合做 FULL OUTER JOIN,用 null-safe 比较(NOT (a <=> b))找出差异字段,快速把"什么时候、哪些行、哪些列"三个问题一次回答掉。
故事 3:8 小时时间偏差——28800 秒的 bug(展示:数据 sense + 细节控)
主叙述:
"迁移验证时发现新旧链路的 trading_duration 有一批记录恰好差 28800 秒。看到 28800 我第一反应就是 8 小时 = UTC+8,方向立刻锁定时区。Root cause 是 TIMESTAMP_NTZ 和 TIMESTAMP 两种类型在 UTC+8 会话下的语义不对称:NTZ 是'无时区裸时间',TIMESTAMP 是'带会话时区解释的时间',混用时同一个字面值被隐式偏移了 8 小时。修复是写入前统一 to_utc_timestamp 显式转换,并把 session timezone 全局定为 UTC 写进规范。"
追问点:Flink SQL 里同样有 TIMESTAMP vs TIMESTAMP_LTZ 的坑,table.local-time-zone 影响窗口边界(天级窗口按哪个时区切)——能接上这句说明你真懂,不是背故事。
故事 4:高峰期写入 hang,定位到分布式共识层(展示:对底层原理的理解深度)
主叙述:
"旧链路有个 INSERT ... ON DUPLICATE KEY UPDATE 任务只在交易高峰期 hang。直觉会怀疑查询本身或锁,但我逐步排除后定位到:卡的是事务 commit 阶段,根因是分布式数据库的 Paxos 日志(clog)复制在高峰期出现 I/O 瓶颈——不是算不动,是'确认写成功'这一步在等多数派。缓解方案三件套:异步提交 + 轮询结果、幂等重试(冲突键更新天然幂等)、大批量任务错峰。"
迁移话术(主动说):"这个经历让我对 exactly-once sink 的 2PC 特别警惕——commit 阶段依赖外部系统的确认,checkpoint 的 notifyCheckpointComplete 卡在 sink commit 上时,表现和这个案例一模一样。所以我做 sink 设计时一定给 commit 阶段单独埋点监控。"
故事 5:回看窗口 90 天→30 天的数据驱动优化(展示:用数据定参数,不拍脑袋)
主叙述: "平仓记录要回查开仓信息,原实现回看 90 天,又慢又重。我先统计了持仓生命周期分布:99.97% 的持仓在 30 天内平仓。于是把回看窗口压到 30 天,剩下 0.03% 走点查兜底(按 position_id 精确回查,命中率 100%),整体把一个依赖物化视图刷新(刷新周期 56 分钟到 3 小时,和 10 分钟调度严重错配)的脆弱链路,改成了确定性的点查链路。"
迁移话术:"这正是 Flink 里定 interval join 区间和 state TTL 的方法论——先统计业务上的真实关联跨度分布,按分位数定窗口,长尾走旁路补偿,而不是拍一个大数让 state 无限膨胀。"
故事 6:三层数据质量验证体系(直接回答 JD 第 5 条"建立科学数据报警体系")
主叙述:
"迁移期间我建了一套 27 个用例的分层验证体系:L1 行级(行数、主键唯一性、分区连续性——用 SEQUENCE + EXPLODE 生成期望分区序列和实际做差集,自动发现掉分区);L2 字段级(新旧链路 FULL OUTER JOIN 逐字段比对,用 null-safe 的 <=> 避免 NULL 比较陷阱);L3 业务级(金额守恒、状态机合法性这类业务不变量)。产出是 HTML 报告,红黄绿一目了然。"
升华(必说):"如果到实时链路,我会把同样的分层思想做成常驻监控:任务层(checkpoint 成功率/时长、消费 lag、反压时间占比、重启次数,Prometheus + Grafana + 分级告警)、数据层(空值率、主键冲突率、流量同环比波动,做成旁路 Flink 作业或 Doris 定时校验)、业务层(核心指标突变检测 + 实时和离线 T+1 对账)。报警要分级——P0 打电话,P2 进群,否则狼来了没人看。"
故事 7:NULL key 伪装成超级热点(小故事,展示 data sense,可在聊倾斜时顺手讲)
"做重复检测时发现某个 key 的重复计数大得离谱,像超级热点。查下去发现是 GROUP BY 把所有 NULL 归为一组,几百万条 NULL key 记录伪装成了'一个超大重复组'。教训:看到倾斜/超大 key,第一步先看是不是 NULL 或默认值(空串、0、-1)聚集,加 IS NOT NULL 或前置清洗再下结论。顺带,当时全量自关联查重复直接 OOM,我改成 CTAS 先物化中间结果再分析——大查询拆阶段物化也是排查时保命的习惯。"
附:批经验 → 实时经验的映射表(自查用)
| 你已有的 | 面试映射到 |
|---|---|
| Paimon snapshot / tag / changelog 增量 | Flink CDC 语义、流式写湖、changelog-producer 模式 |
| 时区 NTZ 事故 | TIMESTAMP vs TIMESTAMP_LTZ、窗口时区边界 |
| 90→30 天回看窗口统计 | interval join 区间、state TTL 的数据驱动设定 |
| 全量覆盖竞态 | 维表一致性、temporal join |
| Paxos commit 瓶颈 | 2PC sink commit 阶段风险与监控 |
| 27 用例分层验证 | 任务/数据/业务三层实时监控报警体系 |
| MT 服务器时间权威 | 游戏客户端时间不可信,event time 取服务端时间 |
| 12.8 亿行、零 NULL 校验 | 大数据量 + 金融级准确性,对标他们月 2 亿场对局 |
附:布隆过滤器速记卡(考前 5 分钟看)
- 六大问题:假阳性(去重=丢数)、不可删(Counting BF ×4 空间)、容量固定(Scalable BF 链)、不可遍历/取数、k 次哈希成本、合并需同参。
- 公式:
m/n = -lnp/(ln2)²,k = (m/n)·ln2;实例:1 亿 key @1% → 114MB、7 个哈希。 - 替代:Cuckoo(可删、低 FPP 更省)、HLL(只算基数不判重!)、Roaring Bitmap(精确、需整数 key)、state+TTL(精确)。
- 出现在:RocksDB SST、HBase/Doris 索引、Redis 防穿透、两级去重(BF 挡"肯定没有"+精确层确认)。
- 一票否决句:"业务不容忍丢数据就不能单独用 BF 去重。"