引言
大多数图数据库的用法是「批量导入 + 在线查询」:夜里把数据灌进去,白天查询。但越来越多的场景要求图「跟着业务实时长出来」——风控要在转账发生的 200 毫秒内判断这笔钱是否流入了一个可疑环,推荐要在用户点击的瞬间更新共现关系,反洗钱要持续维护一张「资金流动图」而不是每天重建一次。流式图处理解决的就是「图数据持续变化时,如何低延迟地写入、更新、计算」的问题。它比流式处理普通表更复杂:图的一条边变化会影响它两端节点的所有邻居,增量计算要维护「局部状态」,而乱序、重复、迟到的事件会让图结构反复抖动。本文讲流式图处理的完整链路:先讲问题定义与实时性需求,再讲 CDC 入图(怎么把关系库的变更映射成图)、增量更新与幂等写入、窗口化子图与状态管理、实时欺诈检测、实时推荐、流式图算法(增量 PageRank 与社区)、与批处理的对照与协同、技术选型与架构,最后是运维监控与常见坑。目标:你能把「实时图」从架构图落到可运行的流水线。
目录
- 1. 流式图处理的问题定义
- 2. CDC 入图:从关系库到图
- 3. 增量更新与幂等写入
- 4. 窗口化子图与状态管理
- 5. 实时欺诈检测
- 6. 实时推荐
- 7. 流式图算法:增量计算
- 8. 与批处理的对照与协同
- 9. 技术选型与架构
- 10. 运维、监控与常见坑
- 速查表
- 延伸阅读
1. 流式图处理的问题定义
为什么图需要「流式」:
- 风控:转账发生时就要判断,不能等 T+1 批处理
- 推荐:用户点击后立即影响下一次推荐
- 反洗钱:资金图持续变化,环可能几分钟内形成又消失
- 供应链/运维:拓扑变化要实时反映到影响面分析
→ 图的「新鲜度」本身就是业务价值
流式图处理与流式表处理的差异:
表处理:一行变化只影响这一行(局部)
图处理:一条边变化影响两端节点的所有邻居(邻域)
- 新增边 → 两端度数变化 → 影响中心性/社区
- 删除节点 → 悬挂边处理 → 图结构一致性
→ 图的变化是「邻域级」的,增量计算要维护邻域状态
延迟与一致性的取舍:
| 模式 | 延迟 | 一致性 | 适用 |
|---|---|---|---|
| 批量重建 | 小时/天 | 强一致(快照) | 离线分析 |
| 微批(Mini-batch) | 秒级 | 最终一致 | 准实时风控 |
| 真流式 | 毫秒/亚秒 | 最终一致 + 幂等 | 实时决策 |
| 混合(Lambda) | 秒级 + 离线校正 | 最终一致 | 高准确要求 |
流式图的三类工作负载:
1. 写入型:把变更事件落到图(CDC 入图)
2. 查询型:对「当前图状态」做低延迟查询(实时风控)
3. 计算型:持续维护图指标(增量 PageRank / 社区)
→ 三类负载的资源诉求不同,常需分开部署
「实时」到底要多实时:
决策类(风控)P99 < 500ms / 推荐类 P99 < 200ms / 监控与图谱维护秒级到分钟级
→ 先定义每个场景的延迟预算,再决定架构
心智:流式图处理的价值是「图的新鲜度」,与流式表处理的关键差异是「变化是邻域级的」——一条边变化影响两端节点的所有邻居,增量计算必须维护邻域状态;按延迟与一致性分模式(批量重建/微批/真流式/混合),并按「写入型、查询型、计算型」三类负载分开部署;先定每个场景的延迟预算,再选架构。
2. CDC 入图:从关系库到图
CDC 的价值:
- 不侵入业务库:不改造写入逻辑,靠读日志捕获变更
- 低延迟:日志一产生就捕获(秒级甚至亚秒)
- 完整:insert / update / delete 都能捕获
→ CDC 是「业务库 → 图」最主流的入图方式
CDC 的三种实现:
1. 日志型(推荐):解析数据库 binlog / WAL
- 优点:低延迟、低侵入、能捕获删除
2. 触发器型:在表上加触发器写变更表
- 优点:实现简单;缺点:侵入业务、有性能损耗
3. 轮询型:定期扫描 updated_at 字段
- 优点:无需数据库权限;缺点:延迟高、捕获不到删除
→ 生产首选日志型 CDC
变更事件的图映射:
关系表 → 图模型的映射规则:
用户表 user → (:User {id, name, ...})
转账表 transfer → (:Account)-[:TRANSFER {amount, at}]->(:Account)
好友表 friendship → (:User)-[:FRIEND]->(:User)
设备登录表 login → (:User)-[:USED_DEVICE {at}]->(:Device)
→ 映射规则要「声明式」配置,而非硬编码在代码里
变更事件的统一格式:
{
"op": "c", "table": "transfer", "ts_ms": 1759300000123,
"after": {"id": "t_90001", "from_account": "a_1001", "to_account": "a_2002",
"amount": 85000.00, "at": "2026-10-01T09:12:03Z"},
"before": null
}
把变更写成 Cypher(幂等 MERGE):
// 入边事件:MERGE 保证幂等(重复消费不产生重复边)
MERGE (a:Account {id: $fromAccount})
MERGE (b:Account {id: $toAccount})
MERGE (a)-[r:TRANSFER {txId: $txId}]->(b)
SET r.amount = $amount, r.at = datetime($at)
删除事件的映射:
// 删除一条关系(幂等:不存在也不报错)
MATCH (a:Account {id: $fromAccount})-[r:TRANSFER {txId: $txId}]->(b:Account)
DELETE r
// 删除节点(先删关系,避免悬挂)
MATCH (n:User {id: $userId})
DETACH DELETE n
CDC 入图的三个难点:
1. 顺序性:同实体变更必须按顺序应用 → 用主键哈希分区保证同实体同分区
2. 删除语义:物理删除传播为图删除,软删除打标 → 图侧区分「真删」与「标记删」
3. Schema 漂移:业务表加字段 → 映射配置要版本化,变更走发布流程
心智:CDC 是「业务库 → 图」的主流入图方式,三种实现里日志型(binlog/WAL)最优(低延迟、低侵入、能捕获删除),触发器与轮询各有硬伤;映射规则要声明式配置(表 → 标签/关系),变更事件用统一格式承载;写图用 MERGE 保证幂等、用主键哈希分区保证同实体顺序、删除要区分物理删与软删除、Schema 漂移要走映射配置的版本化发布。
3. 增量更新与幂等写入
为什么幂等是流式图的生命线:
- 消息系统至少一次投递 → 重复消费必然发生
- 重试、重放、故障恢复都会重复投递
- 非幂等写入 = 重复边、重复计数、图被「污染」
→ 幂等的实现方式决定图能否长期正确
幂等写入的三种手段:
手段 1:MERGE + 业务唯一键(最常用)
MERGE (a)-[r:TRANSFER {txId: $txId}]->(b)
→ 同一 txId 重复写入只产生一条边
手段 2:版本号/时间戳比对(防旧覆盖新)
WHERE r.version < $version SET r.version = $version
手段 3:去重表 / 幂等键缓存(在应用层挡掉重复)
→ 关系用唯一键 MERGE,属性用版本号,应用层再去重一层
乱序事件的处理:
- 事件可能迟到(网络抖动、重试)
- 迟到的「旧」事件不能覆盖「新」状态
- 手段:版本号(单调递增)或事件时间戳比对
- 更复杂:允许「迟到窗口」内的乱序,超窗则进死信
→ 图侧永远「只接受更新的版本」
// 版本号保护:只接受更新的状态
MATCH (a:Account {id: $fromAccount})-[r:TRANSFER {txId: $txId}]->(b:Account)
WHERE r.version IS NULL OR r.version < $version
SET r.status = $status, r.version = $version
Exactly-Once 的真相与写入批量化:
端到端 Exactly-Once 需「消费位点 + 写入」同事务,图库通常不支持
→ 务实做法:至少一次投递 + 幂等写入 = 效果等价
// 攒批后一次写多行:批越大吞吐越高、延迟越高
// 风控类延迟敏感用小批,图谱维护类用大批
UNWIND $events AS e
MERGE (a:Account {id: e.fromAccount})
MERGE (b:Account {id: e.toAccount})
MERGE (a)-[r:TRANSFER {txId: e.txId}]->(b)
SET r.amount = e.amount, r.at = datetime(e.at)
写入压力与背压:
- 上游突发流量 → 图库写入被打满
- 手段:消息队列缓冲 + 消费限流 + 背压反馈
- 监控:写入队列深度、写入 P99、失败率
→ 背压是「保护图库不被冲垮」的必要机制
心智:流式图的生命线是幂等——消息系统至少一次投递决定了重复必然发生;手段是「关系用业务唯一键 MERGE、属性用版本号比对、应用层再去重一层」;乱序事件靠版本号/事件时间戳保证「只接受更新的版本」,超窗进死信;端到端 Exactly-Once 在图库上通常不可得,务实做法是「至少一次投递 + 幂等写入」;写入要攒批(UNWIND)、要有队列缓冲与背压保护图库。
4. 窗口化子图与状态管理
为什么要窗口化:
- 实时计算通常只关心「最近一段时间」的子图
- 例:近 1 小时的转账、近 10 分钟的登录
- 全图计算代价高且大部分是历史噪声
→ 窗口 = 把计算范围从「全图」缩到「热子图」
三种窗口:
滑动窗口(Sliding):近 N 分钟,每步滑动
- 适合:持续监控(近 5 分钟异常登录)
滚动窗口(Tumbling):固定不重叠的 N 分钟桶
- 适合:周期性统计(每 5 分钟聚合)
会话窗口(Session):按活动间隙切分
- 适合:用户会话(30 分钟无操作则结束会话)
→ 按业务语义选窗口,而不是默认滑动
窗口化子图的维护与索引支撑:
// 查询「近 1 小时」的子图
MATCH (a:Account)-[r:TRANSFER]->(b:Account)
WHERE r.at > datetime() - duration('PT1H')
RETURN a.id, b.id, r.amount
// 时间属性必须有索引,否则窗口查询全扫
CREATE INDEX transfer_at_idx FOR ()-[r:TRANSFER]-() ON (r.at)
CREATE INDEX login_at_idx FOR ()-[r:LOGIN]-() ON (r.at)
状态后端的选择:
- 计算框架(Flink)状态:适合「算子内部状态」(计数、聚合)
- 图数据库:适合「需要被查询的图结构」(风控要查图)
- 缓存(Redis):适合「最近窗口的热数据」(低延迟)
→ 三者分工:状态在框架、结构在图库、热数据在缓存
状态规模与清理:
- 窗口状态不清 → 内存持续增长(最终 OOM)
- 清理策略:窗口过期即清、TTL 自动淘汰
- 图库侧:过期边可保留(历史价值)但打标「冷」
→ 「窗口有进有出」是流式图稳定的关键
心智:窗口化把计算范围从全图缩到热子图,三种窗口(滑动/滚动/会话)按业务语义选;时间属性必须有索引否则窗口查询全扫;状态分工是「算子状态在计算框架、图结构在图库、热数据在缓存」;窗口必须有进有出(过期即清、TTL 淘汰),否则状态持续膨胀直至 OOM;恢复时计算框架 checkpoint 与图库恢复点可能不一致,要靠幂等重放对齐。
5. 实时欺诈检测
实时风控的图模式:
1. 环形转账:A→B→C→A 短时间内闭环
2. 资金归集:多个账户短时间转入同一账户
3. 设备聚集:多账户共用同一设备/IP
4. 中介链路:资金经多层「中介账户」快速流转
→ 这些都是「邻域级」模式,适合流式检测
实时环形检测:
// 检测 4 跳内的资金环(限窗口);代价高,实时路径要限深度与窗口
MATCH p = (a:Account {id: $accountId})-[:TRANSFER*2..4]->(a)
WHERE all(r IN relationships(p) WHERE r.at > datetime() - duration('PT1H'))
RETURN p, length(p) AS 环长度
LIMIT 10
更稳的做法:先做「增量标记」把疑似环的账户打标,再由准实时任务确认
增量维护风险分数:
// 收到一笔转账后,更新相关账户的风险计数
MERGE (a:Account {id: $fromAccount})
MERGE (b:Account {id: $toAccount})
MERGE (a)-[r:TRANSFER {txId: $txId}]->(b)
SET r.amount = $amount, r.at = datetime($at)
WITH a, b
SET a.recentOut = coalesce(a.recentOut, 0) + 1,
b.recentIn = coalesce(b.recentIn, 0) + 1
两级判定架构(推荐):
第一级(毫秒级,规则/图查询):命中明显模式 → 直接拦截(黑名单、明显环)
第二级(秒级,图算法/模型):对可疑样本做更深图分析(社区、中心性、GNN)
第三级(离线,全量):批量重算校正误判、发现新模式、回流规则
→ 实时拦截靠第一级,准确性靠后两级
误判与申诉:
- 实时拦截必然有误判 → 要有申诉与人工复核
- 记录「决策依据」(命中了哪条规则/哪个模式)
- 误判样本回流,用于调阈值与规则
→ 可解释性是风控合规的硬要求
心智:实时欺诈检测的核心是「邻域级模式 + 两级判定」——第一级毫秒级规则/图查询直接拦截,第二级秒级图算法/GNN 复核可疑样本,第三级离线批量重算校正并回流规则;实时特征必须「毫秒可算」(窗口入出度、与黑名单限深距离、金额 Z-Score、邻居风险聚合),复杂特征留给第二级;实时拦截必有误判,必须记录决策依据、提供申诉与人工复核、让误判样本回流。
6. 实时推荐
实时推荐的图信号:
1. 实时共现:用户刚点击/购买的商品 → 与哪些商品共现
2. 会话推荐:本次会话内的行为序列 → 下一步推荐
3. 社交影响:好友最近的行为 → 影响推荐
4. 实时热度:窗口内的交互计数 → 提升新内容曝光
→ 实时信号弥补了离线模型的「滞后」
实时共现的更新:
// 用户购买 → 更新商品共现边(滑动窗口内)
MERGE (u:User {id: $userId})
MERGE (p:Product {id: $productId})
MERGE (u)-[:BUY {at: datetime($at)}]->(p)
WITH u, p
MATCH (u)-[:BUY]->(other:Product) WHERE other.id <> p.id
MERGE (p)-[c:CO_OCCUR]->(other)
SET c.recent = coalesce(c.recent, 0) + 1, c.updatedAt = datetime()
实时推荐的查询路径:
// 基于「相似用户最近行为」的实时召回
MATCH (me:User {id: $userId})-[:BUY]->(p:Product)<-[:BUY]-(other:User)
WHERE other.id <> me.id
MATCH (other)-[b:BUY]->(rec:Product)
WHERE NOT (me)-[:BUY]->(rec)
AND b.at > datetime() - duration('PT6H')
RETURN rec.id, count(*) AS score
ORDER BY score DESC
LIMIT 20
实时与离线的融合:
- 离线:协同过滤 / 图嵌入 / GNN 产出候选集(覆盖广)
- 实时:窗口共现 / 会话行为 调整排序(时效强)
- 融合:离线召回 → 实时重排(rerank)
→ 离线管「广度」,实时管「新鲜度」
冷启动与性能约束:
冷启动:新用户用热门/地域/上下文;新商品用内容相似(属性/类目图)
→ 图对冷启动有天然优势(无交互也能靠属性关系找相似)
性能:召回 + 重排总预算 P99 < 200ms,查询深度 ≤ 2、必走索引、强制 LIMIT
→ 实时推荐的最大敌人是超节点(要预计算或缓存)
心智:实时推荐用图信号补足离线模型的滞后——实时共现(滑动窗口内更新 CO_OCCUR)、会话行为、社交影响、窗口热度;架构是「离线召回(广度)+ 实时重排(新鲜度)」;图对冷启动有天然优势(无交互也能靠属性/类目关系找相似);实时路径 P99 预算 200ms 内,查询深度 ≤ 2、必走索引、强制 LIMIT,最大敌人是超节点(要预计算或缓存)。
7. 流式图算法:增量计算
为什么不能每次全量跑算法:
- PageRank / 社区检测在全图上跑一次可能几分钟到几小时
- 但图只变了 0.1% → 全量重算浪费 99.9% 算力
- 增量算法:只更新「受影响的部分」
→ 增量 = 用「局部重算」逼近「全量结果」
增量 PageRank 的思路:
- PageRank 的分数取决于邻居分数
- 新增一条边 → 只影响两端节点及其邻域的分数
- 增量做法:
1) 标记受影响的节点集合(如 3 跳内)
2) 只在这些节点上迭代若干轮
3) 其余节点沿用旧分数
→ 增量结果是近似值,需定期全量校正
增量社区检测的思路:
- 社区结构对局部变化不敏感:检测「跨社区的新边」→ 只在受影响社区内重算
→ 「局部重算 + 定期校正」是增量的通用范式
流式图计算的两种部署:
部署 A:图库内计算(GDS 增量/定期重算)
- 优点:数据不搬家、一致性好
- 缺点:计算与查询争资源
部署 B:流式计算框架内计算(Flink 图算子)
- 优点:吞吐高、与查询隔离
- 缺点:要在框架内维护图状态,复杂度高
→ 中小规模用 A,超大规模用 B
增量结果的校正:
- 增量必然有误差累积(近似算法 + 局部重算)
- 校正手段:周期性全量重算(如每天一次)
- 或「触发式全量」:累积变化超阈值就全量
- 校正期间结果要标记「校正中」
→ 增量 + 定期校正 = 可接受的准确性与成本平衡
结果缓存与物化:
- 增量算出的分数写入节点属性(如 pagerank, communityId)
- 查询直接读属性(O(1)),不再现算
- 更新属性用批量 SET(避免逐条事务)
→ 「算法结果物化为属性」是流式图的标准做法
心智:流式图算法用「局部重算 + 定期校正」替代全量重算——增量 PageRank 只更新受影响邻域、增量社区只重算受影响社区,结果是近似值,靠周期性或触发式全量校正控制误差;部署有「图库内计算(数据不搬家但争资源)」与「流式框架内计算(吞吐高但维护图状态复杂)」两种,中小规模用前者;算法结果物化为节点属性(O(1) 查询),更新用批量 SET。
8. 与批处理的对照与协同
批 vs 流的对照表:
| 维度 | 批处理 | 流处理 |
|---|---|---|
| 延迟 | 分钟到小时 | 毫秒到秒 |
| 数据范围 | 全量快照 | 窗口增量 |
| 一致性 | 强(同一快照) | 最终一致 |
| 算法 | 精确(可多次迭代) | 近似(局部迭代) |
| 复杂度 | 低 | 高 |
| 成本 | 低(按批) | 高(常驻资源) |
| 适用 | 离线分析、报表、全量算法 | 实时决策、实时特征 |
Lambda 架构在图上的落地:
批层:每日全量重算(PageRank、社区、嵌入)→ 写入图属性
流层:实时更新(窗口共现、风险计数)→ 写入图属性
服务层:查询时「批结果 + 流结果」融合(如加权合并)
→ 批管准确性,流管新鲜度
Kappa 架构的取舍:
- 只保留流层,批用「重放历史流」实现
- 优点:一套代码;缺点:重放全量历史成本高
- 图场景常不适用(全量重算图算法用流重放代价太大)
→ 图场景多采用「批 + 流」混合而非纯 Kappa
协同的关键:结果如何合并:
// 服务层融合:批分(稳定)与流分(新鲜)加权
MATCH (u:User {id: $userId})
RETURN 0.7 * coalesce(u.batchScore, 0) + 0.3 * coalesce(u.streamScore, 0) AS score
批流对齐的难点:
- 批结果与流结果时间基准不同(批是快照,流是事件时间),合并要统一口径
- 批重算会「覆盖」流的最新结果 → 需保留流侧增量
→ 合并逻辑要明确「谁覆盖谁、按什么时间」
什么时候不该上流式:
- 业务对延迟不敏感(T+1 报表足够)→ 纯批更省
- 图规模小、变更少 → 定期重建足够
- 团队缺乏流式运维能力 → 复杂度会反噬
→ 流式不是「更先进」,是「更贵」,按需选择
心智:批处理强在精确与低成本、流处理强在低延迟与新鲜度,图上常用「批 + 流」混合(Lambda):批层每日全量重算算法并写属性、流层实时更新窗口特征并写属性、服务层按时间口径加权融合;纯 Kappa 在图场景常不适用(全量图算法重放历史代价太大);批流对齐的难点是时间基准不同与「批覆盖流」,合并逻辑要明确谁覆盖谁;延迟不敏感或规模小的场景,纯批更划算——流式是「更贵」不是「更先进」。
9. 技术选型与架构
技术栈的组合:
消息层:Kafka / Pulsar(缓冲、分区、顺序保证)
变更捕获:Debezium(日志型 CDC)
流计算:Kafka Streams(轻) / Flink(重、状态强)
图存储:Neo4j / JanusGraph / NebulaGraph;图算法:GDS / 自研
缓存:Redis(热窗口数据)→ 按规模与团队能力组合,不追求全家桶
架构参考(实时风控):
业务库 → Debezium CDC → Kafka(按账户 ID 分区)
↓
流计算(窗口聚合 + 规则)
↓
图库写入(MERGE 幂等)→ 实时图查询
↓
决策服务(毫秒级返回)
分区键的选择:
- 必须保证「同一实体的变更有序」
- 风控:按账户 ID 分区(同账户的转账有序)
- 推荐:按用户 ID 分区(同用户行为有序)
- 错误示范:按事件时间分区(同实体可能落在不同分区)
→ 分区键 = 「顺序敏感的最小实体」
读写分离与资源隔离:
- 写入(CDC 入图)与查询(风控决策)争资源
- 方案:读走只读副本、写走主库;或双集群
- 图算法计算单独部署(避免抢查询资源)
→ 三类负载(写/查/算)尽量隔离
容量估算与选型的三个问题:
1. 延迟预算多少(决定真流式还是微批)
2. 状态规模多大(决定 Flink 还是图库内计算)
3. 团队会不会运维(决定复杂度上限)
→ 三个问题回答完,选型基本确定
心智:技术栈按「消息(Kafka)+ CDC(Debezium)+ 流计算(Kafka Streams 轻 / Flink 重)+ 图库 + 算法(GDS 或自研)+ 缓存」组合;分区键必须是「顺序敏感的最小实体」(账户/用户 ID),绝不能按事件时间分区;三类负载(写/查/算)要隔离(只读副本、独立算法集群);容量估算先估写入压力与窗口状态,再估图库内存;选型回答三个问题——延迟预算、状态规模、团队运维能力。
10. 运维、监控与常见坑
关键监控指标:
1. 端到端延迟(事件产生 → 图可见)
2. 消费 Lag(消息堆积)
3. 图写入 QPS / P99 / 失败率
4. 幂等冲突率(MERGE 命中已有键的比例)
5. 窗口状态大小(是否持续增长)
6. 增量算法误差(与全量结果的偏差)
→ 前三个是「健康度」,后三个是「正确性」
常见坑清单:
坑 1:非幂等写入 → 重复消费产生重复边/重复计数
坑 2:按事件时间分区 → 同实体变更乱序,状态错乱
坑 3:无版本号保护 → 迟到事件覆盖新状态
坑 4:窗口状态不清理 → 内存持续增长直至 OOM
坑 5:时间属性无索引 → 窗口查询全表扫描
坑 6:删除事件未处理 → 图里留着已删除的实体
坑 7:实时路径上跑深度遍历 → 超节点拖垮延迟
坑 8:增量算法不校正 → 误差累积,结果越跑越偏
坑 9:批流结果直接相加 → 时间口径不同,结果失真
坑 10:CDC 与图写入无背压 → 上游突发冲垮图库
坑 11:把「框架 EOS」当「业务幂等」→ 端到端仍可能重复
坑 12:流式与在线查询共用实例 → 互相争资源
上线检查清单:
[ ] 分区键保证同实体有序
[ ] 写入全链路幂等(唯一键 MERGE + 版本号)
[ ] 窗口状态有 TTL 与清理策略
[ ] 时间属性建索引
[ ] 删除事件有处理路径
[ ] 实时路径限深度、强制 LIMIT
[ ] 增量算法有定期全量校正
[ ] 批流合并有明确时间口径
[ ] 背压与限流机制就绪
[ ] 三类负载资源隔离
故障恢复演练:
消费位点回退重放(验幂等)/ 图库恢复 + 流重放(验状态一致)
上游突发(验背压)/ 算法校正失败(验降级)→ 覆盖重复、丢失、突发、失败四类
心智:流式图的运维监控看六项——端到端延迟、消费 Lag、写入失败率(健康度)+ 幂等冲突率、窗口状态大小、增量误差(正确性);十二个坑集中在四处——幂等缺失(重复边/乱序覆盖)、状态失控(窗口不清理 OOM)、查询失控(无索引/深遍历/超节点)、架构缺陷(批流口径错、无背压、负载不隔离、误信框架 EOS);上线前过检查清单,演练覆盖重复、丢失、突发、失败四类故障。
速查表
窗口与状态速记:
| 主题 | 结论 |
|---|---|
| 幂等 | 唯一键 MERGE + 版本号 + 应用层去重 |
| 顺序 | 分区键 = 顺序敏感的最小实体 |
| 乱序 | 版本号/事件时间,只接受更新的版本 |
| 窗口 | 滑动/滚动/会话,按业务语义选 |
| 状态清理 | TTL 淘汰,窗口有进有出 |
| 状态分工 | 算子在框架、结构在图库、热数据在缓存 |
| 增量算法 | 局部重算 + 定期全量校正 |
| 批流协同 | 批管准确、流管新鲜,服务层按时口径融合 |
| 背压 | 队列缓冲 + 消费限流 + 失败降级 |
| 隔离 | 写/查/算三类负载分开部署 |
延迟预算参考:
风控决策 P99 < 500ms / 推荐 P99 < 200ms
监控 秒级 / 图谱维护 秒级到分钟级
一句话记忆:流式图处理的价值是「图的新鲜度」,与流式表处理的关键差异是「变化是邻域级的」——一条边变化影响两端所有邻居,增量计算必须维护邻域状态;入图主流是日志型 CDC(Debezium + Kafka),映射规则声明式配置;生命线是幂等(关系用业务唯一键 MERGE、属性用版本号、应用层再去重),乱序靠版本号保证「只接受更新的版本」,端到端 Exactly-Once 在图库上通常不可得,务实做法是「至少一次投递 + 幂等写入」;窗口化把计算范围从全图缩到热子图(滑动/滚动/会话按业务选),时间属性必须建索引,窗口必须「有进有出」否则状态膨胀 OOM,状态分工是「算子在框架、结构在图库、热数据在缓存」;实时风控靠「两级判定」(毫秒级规则拦截 + 秒级图算法复核 + 离线校正回流),实时推荐靠「离线召回 + 实时重排」;流式图算法用「局部重算 + 定期全量校正」替代全量重算并把结果物化为节点属性;批流协同是 Lambda——批管准确、流管新鲜、服务层按时口径融合,纯 Kappa 在图场景常不适用;分区键必须是「顺序敏感的最小实体」,写/查/算三类负载要隔离;运维看端到端延迟、消费 Lag、写入失败率与幂等冲突率、窗口状态大小、增量误差六项,演练覆盖重复、丢失、突发、失败四类——把幂等、窗口、背压、隔离四件事做扎实,实时图才稳得住。
延伸阅读
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。