使用Java、Kafka Streams和向量相似度实现实时欺诈检测
本文基于我和同事 Tim Kelly 在 DevNexus 2026 上发表的一次演讲,该会议在亚特兰大举办,是世界上最大的 Java 会议之一。
想象你正在超市购物或在线购物。你刷了卡,或点击了“立即购买”。就在那一刻,你的银行必须做出决定:这笔交易是合法的,还是欺诈?而且它必须在几毫秒内完成,并且判断正确。如果它拦截了一笔真实交易,会让客户感到沮丧。如果它批准了一笔欺诈交易,就会导致财务损失。那么问题来了:银行究竟是如何做到这一点的?幕后发生了什么?在本文中,我们将探讨一种真实世界的解决方案,使用 Java、Apache Kafka、Kafka Streams 以及 MongoDB 的向量相似度来解决这个问题。
我们要构建什么
从高层来看,支付系统的工作方式如下:创建一笔交易,经过验证步骤,然后做出最终决定:批准或拒绝。
现在,让我们聚焦于验证步骤,因为欺诈检测实际上就发生在这里。在我们的方案中,这个验证不是单一检查,而是由多个阶段协同工作的管道:
- 防护栏: 基于规则的检查,立即捕获明显且已知的欺诈模式,例如不可能旅行或不寻常的交易速度。
- 向量评分: AI 驱动的相似度匹配,将交易与已知模式进行比较,以检测更细微或行为层面的欺诈。
- 决策: 最终裁决基于两层的组合输出:将交易路由到
approved_transactions集合,或路由到suspicious_transactions集合,以便进一步审查或明确拦截。
决策流程相当直接。如果一笔交易在任一环节被标记,无论是被防护栏还是向量评分阶段标记,我们都会将其视为可疑,并存入 suspicious_transactions 集合。然后,可以根据场景对其进行审查、补充更多上下文,甚至拦截。另一方面,如果它顺利通过两层而没有问题,我们就认为它是安全的,并将其存入 approved_transactions 集合。
挑战:毫秒级决策与故障下的韧性
在欺诈检测系统中,尤其是金融应用中,速度不仅重要,而且至关重要。每笔交易都必须实时分析,通常是在高负载下,每秒发生数千个事件。同时,系统不能因为单个组件故障而崩溃。
这带来了两个不可协商的约束:
- 决策必须在毫秒内完成
- 即使单个组件发生故障,系统也必须保持韧性
在我们深入探讨每一层如何工作之前,重要的是理解这些约束如何塑造整个架构。
自然而然的第一反应是使用简单的同步架构来构建。支付服务调用欺诈服务,欺诈服务再查询数据库,然后流程继续到通知服务。这在初期运行良好。它简单且易于理解。
但问题是,一切都紧密连接。一个服务依赖下一个服务。现在想象数据库响应稍慢。这种延迟开始影响整个流程。欺诈服务变慢,支付服务必须等待,通知服务也会被延迟。
在规模扩大后,这会成为真正的问题。一个缓慢的组件会影响其他所有组件。那么,我们如何打破这种依赖,让每个服务更加独立?这就是事件驱动架构发挥作用的地方。
策略:用事件驱动架构解耦系统
与其让服务彼此直接调用,我们引入事件日志作为通信骨干。
Apache Kafka 允许支付服务将交易事件发布到某个主题,而不需要知道谁会消费它。欺诈服务、通知服务以及任何其他感兴趣的组件都可以独立订阅该主题。
这从根本上改变了系统的行为:
- 服务不再紧密耦合
- 每个服务都可以失败或重启,而不影响其他服务
- 系统变得更有韧性,也更容易扩展
至此,我们已经移除了服务之间的直接依赖,并解决了韧性问题。但我们还有另一个挑战。
无需数据库往返的有状态检查
许多欺诈检测技术依赖于理解近期活动。例如,系统可能需要评估一张卡在短时间窗口内是如何被使用的,然后再决定一笔新交易是否可疑。
在传统架构中,这通常需要为每笔交易查询数据库。在高吞吐量下,这些往返会增加延迟,很快成为瓶颈。在一个必须在毫秒内做出决策的系统中,这种方式无法扩展。
那么问题就变成了:我们如何在不每次都访问数据库的情况下执行这类检查?
用 Kafka Streams 解决
Kafka Streams 是一个轻量级 Java 库,运行在你的应用程序内部,在事件流经 Kafka 主题时处理它们。更重要的是,它允许我们在本地维护状态。
我们不再为每笔交易查询远程数据库,而是将必要数据保留在处理事件的同一个进程内,使用由 RocksDB 支持的嵌入式存储。
在我们的示例中,我们希望分析一张卡在短时间窗口内的交易活动。为此,我们按卡号对事件进行分组。每个卡号成为一个键,对于每个键,系统维护一个定义时间窗口内近期活动的持续视图。这意味着我们从关键路径中消除了网络调用和数据库往返,只依赖本地查找。对于我们的用例,这带来了很大差异。
我们避免了不必要的延迟,并让一切快到足以进行实时决策,而不依赖外部系统。你可以直接在代码中看到这一点发生的位置:
这一行是关键:
它告诉 Kafka Streams 使用嵌入式存储将状态保留在本地。随着新交易到达,每张卡的状态就在应用程序中更新,因此系统可以实时跟踪近期活动。
至此,我们可以在不依赖数据库调用的情况下实时评估交易活动。但这引出了一个重要问题:我们究竟在度量什么?哪些模式或信号应该触发欺诈决策?这就是防护栏发挥作用的地方。
第 1 阶段:防护栏,快速的基于规则验证
防护栏是我们流程中的第一层。它们只是基于规则的简单检查,用于尽早捕获最明显的欺诈案例。
在实践中,防护栏可以是任何对你的领域有意义的业务或风险规则。它们完全可定制,并随着新欺诈模式出现而演进。在我们的案例中,我们将实现几个常见示例来说明这种方法。
不可能旅行检查
这条规则 检查同一张卡是否出现在鉴于交易间隔时间而不合理的地点。
例如,如果一张卡在纽约使用,然后不到一小时后出现在圣保罗,这并不现实。显然出了问题,我们将其视为强欺诈信号。
速度检查
速度 检查关注一张卡在短时间内被使用的次数。欺诈中非常常见的一种模式是,有人先通过多笔小额交易测试一张被盗卡,然后再尝试更大金额。
在我们的案例中,如果同一张卡在一分钟内被使用超过三次,我们就标记它并拦截该交易。
这些只是示例。在真实系统中,防护栏可以根据你的用例、风险容忍度和业务规则包含各种各样的检查。
基于规则的检测在哪里不够用
防护栏快速且有效,但它们有盲点。考虑这样一笔交易:
- 金额:$2.00
- 城市:与持卡人常用城市相同
- 卡:和往常一样
- 活动:表面上看没有任何异常
这笔交易会通过所有防护栏检查。但是,如果加油站 $2 的小额扣款是一种已知欺诈模式呢?这是在大额欺诈购买之前常用的经典卡片测试技术。仅靠规则无法捕获这一点,因为系统需要提前知道该模式。这就是需要另一种方法的地方。
第 2 阶段:向量评分,行为欺诈检测
此时,问题变了。我们不再问:“这笔交易是否违反了某条规则?”而是问:“这笔交易的行为是否像我们以前见过的东西?”
为了回答这个问题,我们需要一种基于行为而不是精确值来比较交易的方法。例如:
- 在咖啡店消费 $12
- 在面包店消费 $15
它们并不完全相同,但从行为上看非常相似。另一方面:
- 凌晨 3 点的 $10,000 加密货币交易
这显然非常不同。那么,我们如何以一种系统能够理解的方式来表示这个想法?
从交易到向量
为了比较行为,我们需要一种一致的方式来表示交易。我们不再分别处理金额、商户或位置等原始字段,而是将整笔交易转换为单一的结构化表示。一个简单的理解方式是:
我们取这样一笔交易:
并将其转换为描述性文本:
这段文本随后被传递给嵌入模型,生成一个向量,例如:
[0.0392, 0.9323, -0.0323, ...]
在我们的案例中,我们使用 voyage-finance-2 模型,它专门针对来自 Voyage AI 的金融数据。该模型生成一个 1024 维向量,其中每个值都有助于捕捉交易行为的不同方面。
在应用程序中,这个嵌入生成步骤可以通过 Spring 中的简单 HTTP 客户端 实现。例如:
这个客户端的作用是将交易文本发送到嵌入 API,并接收生成的向量。
该模型将描述映射到一个向量空间,其中相似行为被放置得更近。这意味着行为相似的交易会生成彼此接近的向量。
为了让这一点更直观,你可以想象一个简化的二维空间。在这个空间中:
- 咖啡或杂货等日常消费倾向于聚集在一起
- 相似类型的交易保持接近
- 不寻常的模式,例如深夜加密货币交易,则相距很远
实际上,这个空间不是二维的,而是有数百或数千个维度。这正是该模型能够捕捉原始交易数据中不明显的细微模式的原因。结果是交易的数值表示,它保留了其行为,使得可以使用向量相似度高效地比较交易。
构建基线:为欺诈模式集合播种
现在我们可以将交易转换为向量,还需要一些东西来比较它们。向量相似度只有在有参考数据集时才有效。
这意味着我们需要一个已知模式集合,代表正常和欺诈行为。为此,我们在 MongoDB 中创建一个 fraud_patterns 集合。该集合中的每个文档包含:
- 原始交易数据
- 欺诈标签(true 或 false)
- 对应的嵌入
一个正常模式可能如下所示:
一个欺诈模式可能如下所示:
这个数据集充当比较的基线。我们不再问“这笔交易是否欺诈?”,而是问:“这笔交易看起来像我们已经知道的东西吗?”
这个集合如何演化
这个集合不需要一开始就完美。它通常从一小组精选的已知模式开始。随着系统在生产中运行:
- 可以添加已确认的欺诈案例
- 可以立即引入新模式
- 系统会随着时间推移而改进
这里一个重要优势是,我们不需要重新训练模型。我们只需向数据集添加新示例。
使用 MongoDB 运行向量搜索
要运行向量搜索,我们需要创建一个向量索引。这个索引允许 MongoDB 高效地搜索相似向量。
这里的每个字段都很重要:
- type: "vector" 告诉 MongoDB 该字段包含嵌入
- path: "embedding" 定义向量存储在文档中的位置
- numDimensions: 1024 必须与模型生成的向量大小匹配
- similarity: "cosine" 定义向量之间的相似度如何计算
如果维度数量与嵌入模型不匹配,索引将无法正常工作。
执行搜索
一旦索引就位,MongoDB 就可以高效地找到最近的向量。通过 Spring Data MongoDB,这变得非常简单:
这个方法 将交易嵌入作为输入,并返回集合中存储的最相似模式。
并非所有匹配都相同:引入阈值
运行搜索后,MongoDB 会返回最接近的匹配及其相似度分数。例如:
- 结果 1 = 分数:0.94
- 结果 2 = 分数:0.91
- 结果 3 = 分数:0.72
- 结果 4 = 分数:0.60
每个分数告诉我们一个存储模式与我们正在分析的交易有多相似。在实践中,即使相似度相对较低,向量搜索也会返回可用的最接近匹配。
这就是我们引入阈值的原因。阈值定义了结果被视为相关所需的最低相似度分数。
任何低于该值的结果都会被忽略。选择正确的阈值至关重要:
- 太低 → 许多误报
- 太高 → 漏掉欺诈案例
理想值取决于你的领域和风险容忍度。