Ohhnews

分类导航

$ cd ..
foojay原文

使用Java、Kafka Streams和向量相似度实现实时欺诈检测

#实时欺诈检测#kafka streams#java#向量相似度#mongodb

本文基于我和同事 Tim Kelly 在 DevNexus 2026 上发表的一次演讲,该会议在亚特兰大举办,是世界上最大的 Java 会议之一。

想象你正在超市购物或在线购物。你刷了卡,或点击了“立即购买”。就在那一刻,你的银行必须做出决定:这笔交易是合法的,还是欺诈?而且它必须在几毫秒内完成,并且判断正确。如果它拦截了一笔真实交易,会让客户感到沮丧。如果它批准了一笔欺诈交易,就会导致财务损失。那么问题来了:银行究竟是如何做到这一点的?幕后发生了什么?在本文中,我们将探讨一种真实世界的解决方案,使用 Java、Apache Kafka、Kafka Streams 以及 MongoDB 的向量相似度来解决这个问题。

我们要构建什么

从高层来看,支付系统的工作方式如下:创建一笔交易,经过验证步骤,然后做出最终决定:批准或拒绝。

现在,让我们聚焦于验证步骤,因为欺诈检测实际上就发生在这里。在我们的方案中,这个验证不是单一检查,而是由多个阶段协同工作的管道:

  1. 防护栏: 基于规则的检查,立即捕获明显且已知的欺诈模式,例如不可能旅行或不寻常的交易速度。
  2. 向量评分: AI 驱动的相似度匹配,将交易与已知模式进行比较,以检测更细微或行为层面的欺诈。
  3. 决策: 最终裁决基于两层的组合输出:将交易路由到 approved_transactions 集合,或路由到 suspicious_transactions 集合,以便进一步审查或明确拦截。

决策流程相当直接。如果一笔交易在任一环节被标记,无论是被防护栏还是向量评分阶段标记,我们都会将其视为可疑,并存入 suspicious_transactions 集合。然后,可以根据场景对其进行审查、补充更多上下文,甚至拦截。另一方面,如果它顺利通过两层而没有问题,我们就认为它是安全的,并将其存入 approved_transactions 集合。

挑战:毫秒级决策与故障下的韧性

在欺诈检测系统中,尤其是金融应用中,速度不仅重要,而且至关重要。每笔交易都必须实时分析,通常是在高负载下,每秒发生数千个事件。同时,系统不能因为单个组件故障而崩溃。

这带来了两个不可协商的约束:

  • 决策必须在毫秒内完成
  • 即使单个组件发生故障,系统也必须保持韧性

在我们深入探讨每一层如何工作之前,重要的是理解这些约束如何塑造整个架构。

自然而然的第一反应是使用简单的同步架构来构建。支付服务调用欺诈服务,欺诈服务再查询数据库,然后流程继续到通知服务。这在初期运行良好。它简单且易于理解。

但问题是,一切都紧密连接。一个服务依赖下一个服务。现在想象数据库响应稍慢。这种延迟开始影响整个流程。欺诈服务变慢,支付服务必须等待,通知服务也会被延迟。

在规模扩大后,这会成为真正的问题。一个缓慢的组件会影响其他所有组件。那么,我们如何打破这种依赖,让每个服务更加独立?这就是事件驱动架构发挥作用的地方。

策略:用事件驱动架构解耦系统

与其让服务彼此直接调用,我们引入事件日志作为通信骨干。

Apache Kafka 允许支付服务将交易事件发布到某个主题,而不需要知道谁会消费它。欺诈服务、通知服务以及任何其他感兴趣的组件都可以独立订阅该主题。

这从根本上改变了系统的行为:

  • 服务不再紧密耦合
  • 每个服务都可以失败或重启,而不影响其他服务
  • 系统变得更有韧性,也更容易扩展

至此,我们已经移除了服务之间的直接依赖,并解决了韧性问题。但我们还有另一个挑战。

无需数据库往返的有状态检查

许多欺诈检测技术依赖于理解近期活动。例如,系统可能需要评估一张卡在短时间窗口内是如何被使用的,然后再决定一笔新交易是否可疑。

在传统架构中,这通常需要为每笔交易查询数据库。在高吞吐量下,这些往返会增加延迟,很快成为瓶颈。在一个必须在毫秒内做出决策的系统中,这种方式无法扩展。

那么问题就变成了:我们如何在不每次都访问数据库的情况下执行这类检查?

用 Kafka Streams 解决

Kafka Streams 是一个轻量级 Java 库,运行在你的应用程序内部,在事件流经 Kafka 主题时处理它们。更重要的是,它允许我们在本地维护状态。

我们不再为每笔交易查询远程数据库,而是将必要数据保留在处理事件的同一个进程内,使用由 RocksDB 支持的嵌入式存储。

在我们的示例中,我们希望分析一张卡在短时间窗口内的交易活动。为此,我们按卡号对事件进行分组。每个卡号成为一个键,对于每个键,系统维护一个定义时间窗口内近期活动的持续视图。这意味着我们从关键路径中消除了网络调用和数据库往返,只依赖本地查找。对于我们的用例,这带来了很大差异。

我们避免了不必要的延迟,并让一切快到足以进行实时决策,而不依赖外部系统。你可以直接在代码中看到这一点发生的位置:

$ java
return stream
    .filter((key, tx) -> tx != null)
    .selectKey((key, tx) -> tx.cardNumber())
    .groupByKey(Grouped.with(Serdes.String(), transactionSerde))
    .aggregate(
        FraudDetectionState::empty,
        (cardNumber, newTransaction, currentState) -> { countTransactionsInWindow(); },
        Materialized.as("fraud-detection-store")...
    )
    .toStream()

这一行是关键:

$ java
Materialized.as("fraud-detection-store")

它告诉 Kafka Streams 使用嵌入式存储将状态保留在本地。随着新交易到达,每张卡的状态就在应用程序中更新,因此系统可以实时跟踪近期活动。

至此,我们可以在不依赖数据库调用的情况下实时评估交易活动。但这引出了一个重要问题:我们究竟在度量什么?哪些模式或信号应该触发欺诈决策?这就是防护栏发挥作用的地方。

第 1 阶段:防护栏,快速的基于规则验证

防护栏是我们流程中的第一层。它们只是基于规则的简单检查,用于尽早捕获最明显的欺诈案例。

在实践中,防护栏可以是任何对你的领域有意义的业务或风险规则。它们完全可定制,并随着新欺诈模式出现而演进。在我们的案例中,我们将实现几个常见示例来说明这种方法。

不可能旅行检查

这条规则 检查同一张卡是否出现在鉴于交易间隔时间而不合理的地点。

例如,如果一张卡在纽约使用,然后不到一小时后出现在圣保罗,这并不现实。显然出了问题,我们将其视为强欺诈信号。

$ java
public boolean isImpossibleTravelFraud(Transaction previous, Transaction current) {
    double distanceKm = calculateDistanceKm(
        previous.latitude(), previous.longitude(),
        current.latitude(), current.longitude()
    );
    Duration timeBetween = Duration.between(
        previous.transactionTime(),
        current.transactionTime()
    ).abs();
    // block if physically impossible
}

速度检查

速度 检查关注一张卡在短时间内被使用的次数。欺诈中非常常见的一种模式是,有人先通过多笔小额交易测试一张被盗卡,然后再尝试更大金额。

在我们的案例中,如果同一张卡在一分钟内被使用超过三次,我们就标记它并拦截该交易。

$ java
public boolean isVelocityFraud(long txCountInWindow, long maxAllowedPerWindow) {
    return txCountInWindow > maxAllowedPerWindow;
}

这些只是示例。在真实系统中,防护栏可以根据你的用例、风险容忍度和业务规则包含各种各样的检查。

基于规则的检测在哪里不够用

防护栏快速且有效,但它们有盲点。考虑这样一笔交易:

  • 金额:$2.00
  • 城市:与持卡人常用城市相同
  • 卡:和往常一样
  • 活动:表面上看没有任何异常

这笔交易会通过所有防护栏检查。但是,如果加油站 $2 的小额扣款是一种已知欺诈模式呢?这是在大额欺诈购买之前常用的经典卡片测试技术。仅靠规则无法捕获这一点,因为系统需要提前知道该模式。这就是需要另一种方法的地方。

第 2 阶段:向量评分,行为欺诈检测

此时,问题变了。我们不再问:“这笔交易是否违反了某条规则?”而是问:“这笔交易的行为是否像我们以前见过的东西?”

为了回答这个问题,我们需要一种基于行为而不是精确值来比较交易的方法。例如:

  • 在咖啡店消费 $12
  • 在面包店消费 $15

它们并不完全相同,但从行为上看非常相似。另一方面:

  • 凌晨 3 点的 $10,000 加密货币交易

这显然非常不同。那么,我们如何以一种系统能够理解的方式来表示这个想法?

从交易到向量

为了比较行为,我们需要一种一致的方式来表示交易。我们不再分别处理金额、商户或位置等原始字段,而是将整笔交易转换为单一的结构化表示。一个简单的理解方式是:

我们取这样一笔交易:

$ cat
{ "amount": 15.90, "merchant": "Gas Station", "city": "New York" }

并将其转换为描述性文本:

$ java
String textToEmbed = "merchant=%s, city=%s, amount=%s";

这段文本随后被传递给嵌入模型,生成一个向量,例如:

[0.0392, 0.9323, -0.0323, ...]

在我们的案例中,我们使用 voyage-finance-2 模型,它专门针对来自 Voyage AI 的金融数据。该模型生成一个 1024 维向量,其中每个值都有助于捕捉交易行为的不同方面。

在应用程序中,这个嵌入生成步骤可以通过 Spring 中的简单 HTTP 客户端 实现。例如:

$ java
@HttpExchange(
    url = "/v1/embeddings",
    contentType = MediaType.APPLICATION_JSON_VALUE,
    accept = MediaType.APPLICATION_JSON_VALUE
)
public interface VoyageEmbeddingsClient {
    @PostExchange
    EmbeddingsResponse embed(@RequestBody EmbeddingsRequest body);
}

这个客户端的作用是将交易文本发送到嵌入 API,并接收生成的向量。

该模型将描述映射到一个向量空间,其中相似行为被放置得更近。这意味着行为相似的交易会生成彼此接近的向量。

为了让这一点更直观,你可以想象一个简化的二维空间。在这个空间中:

  • 咖啡或杂货等日常消费倾向于聚集在一起
  • 相似类型的交易保持接近
  • 不寻常的模式,例如深夜加密货币交易,则相距很远

实际上,这个空间不是二维的,而是有数百或数千个维度。这正是该模型能够捕捉原始交易数据中不明显的细微模式的原因。结果是交易的数值表示,它保留了其行为,使得可以使用向量相似度高效地比较交易。

构建基线:为欺诈模式集合播种

现在我们可以将交易转换为向量,还需要一些东西来比较它们。向量相似度只有在有参考数据集时才有效。

这意味着我们需要一个已知模式集合,代表正常和欺诈行为。为此,我们在 MongoDB 中创建一个 fraud_patterns 集合。该集合中的每个文档包含:

  • 原始交易数据
  • 欺诈标签(true 或 false)
  • 对应的嵌入

一个正常模式可能如下所示:

$ cat
{
  "fraud": false,
  "transaction": {
    "amount": 13.45,
    "merchant": "Coffee Shop"
  },
  "embedding": [ ... ]
}

一个欺诈模式可能如下所示:

$ cat
{
  "fraud": true,
  "transaction": {
    "amount": 87422.45,
    "merchant": "Crypto Exchange"
  },
  "embedding": [ ... ]
}

这个数据集充当比较的基线。我们不再问“这笔交易是否欺诈?”,而是问:“这笔交易看起来像我们已经知道的东西吗?”

这个集合如何演化

这个集合不需要一开始就完美。它通常从一小组精选的已知模式开始。随着系统在生产中运行:

  • 可以添加已确认的欺诈案例
  • 可以立即引入新模式
  • 系统会随着时间推移而改进

这里一个重要优势是,我们不需要重新训练模型。我们只需向数据集添加新示例。

使用 MongoDB 运行向量搜索

要运行向量搜索,我们需要创建一个向量索引。这个索引允许 MongoDB 高效地搜索相似向量。

$ cat
{
  "fields": [
    {
      "type": "vector",
      "path": "embedding",
      "numDimensions": 1024,
      "similarity": "cosine"
    }
  ]
}

这里的每个字段都很重要:

  • type: "vector" 告诉 MongoDB 该字段包含嵌入
  • path: "embedding" 定义向量存储在文档中的位置
  • numDimensions: 1024 必须与模型生成的向量大小匹配
  • similarity: "cosine" 定义向量之间的相似度如何计算

如果维度数量与嵌入模型不匹配,索引将无法正常工作。

执行搜索

一旦索引就位,MongoDB 就可以高效地找到最近的向量。通过 Spring Data MongoDB,这变得非常简单:

$ java
@Repository
public interface FraudPatternRepository extends MongoRepository<FraudPattern, String> {

    @VectorSearch(
        indexName = "fraud_patterns_vector_index",
        limit = "10",
        numCandidates = "200"
    )
    SearchResults<FraudPattern> searchTopFraudPatternsByEmbeddingNear(
        Vector vector,
        Score score
    );
}

这个方法 将交易嵌入作为输入,并返回集合中存储的最相似模式。

并非所有匹配都相同:引入阈值

运行搜索后,MongoDB 会返回最接近的匹配及其相似度分数。例如:

  • 结果 1 = 分数:0.94
  • 结果 2 = 分数:0.91
  • 结果 3 = 分数:0.72
  • 结果 4 = 分数:0.60

每个分数告诉我们一个存储模式与我们正在分析的交易有多相似。在实践中,即使相似度相对较低,向量搜索也会返回可用的最接近匹配。

这就是我们引入阈值的原因。阈值定义了结果被视为相关所需的最低相似度分数。

$ config
similarity-threshold: 0.95

任何低于该值的结果都会被忽略。选择正确的阈值至关重要:

  • 太低 → 许多误报
  • 太高 → 漏掉欺诈案例

理想值取决于你的领域和风险容忍度。