编程 Venice 深度拆解:LinkedIn 如何用衍生数据平台承载每日 1.2PB 数据写入与万亿级特征服务

2026-07-29 12:16:37 +0800 CST views 6

Venice 深度拆解:LinkedIn 如何用衍生数据平台承载每日 1.2PB 数据写入与万亿级特征服务

写在前面

当我们讨论大规模数据系统时,MySQL、Redis、Kafka 这些名字几乎成了口头禅。但如果我把 LinkedIn 的真实数字摆出来,你可能会重新思考「大规模」的定义:

  • 每日写入量:1.2 PB,约 3 万亿行
  • 峰值写入吞吐:50 GB/s,113 万行/秒
  • 平均写入吞吐:14 GB/s,39 万行/秒
  • 单服务器写入:200 MB/s,40 万行/秒(同时还要 Serving 读取请求)
  • 峰值读取 QPS:4500 万次 Key 查找/秒
  • 生产集群规模:1800+ 数据集,300+ 应用

这些数字不是数据库厂商的白皮书宣传,而是 LinkedIn 在 2022 年将 Venice 开源时,公开发布的真实生产数据。Venice 是 LinkedIn 自研的「衍生数据平台」(Derived Data Platform),从 2016 年底投产至今,持续扩张,替代了 Voldemort 及多个自研系统,成为 LinkedIn AI 基础设施的核心存储层。

更有意思的是它的技术定位:不是 OLTP 数据库,不是 OLAP 数据仓库,也不是传统 KV 存储,而是一个 专门为「衍生数据」设计的分布式存储系统——数据来源于批处理(Spark/Hadoop)和流处理(Samza),服务于在线推理(Feature Store、推荐、搜索)。

这个定位极其精准。它不试图做通用数据库做的事,而是把「衍生数据从离线世界高效传递到在线世界」这件事做到极致。本文将深入拆解 Venice 的架构设计、读写路径、多活机制,以及它在 Feature Store 场景中的独特价值。


一、背景:为什么 LinkedIn 需要一个新的数据系统

1.1 衍生数据:一个被低估的系统设计难题

在互联网公司中,大部分数据可以分为两类:

主数据(Primary Data):用户发的帖子、订单记录、商品信息——这类数据通常由 OLTP 数据库承载,强一致性写入,直接面向用户。

衍生数据(Derived Data):用户的行为特征、推荐分数、搜索索引、实时统计——这类数据是由主数据经过 ML 模型、聚合计算、ETL 流水线「加工」出来的。它的特点是:

  • 写入异步化:数据来自批处理(天级/小时级)和流处理(秒级),不是直接用户请求触发
  • 读取低延迟:在线推理场景(P99 < 10ms)要求毫秒级响应
  • 数据量巨大:一个特征数据集可能包含数十亿行embedding向量
  • 无需强一致性:衍生数据本身就是「计算结果」,允许一定延迟

LinkedIn 的核心业务——「你可能认识的人」「你可能感兴趣的职位」「Feed 推荐」——本质都是衍生数据。这些场景不需要毫秒级的数据新鲜度(推荐不需要你刚点了一个赞就立刻反映),但需要高吞吐的写入管道低延迟的在线读取

1.2 旧系统的困境:Voldemort 的局限性

LinkedIn 此前的主力 KV 存储是 Voldemort(是的,这是 LinkedIn 的另一款开源产品,名字来自《哈利·波特》中的黑魔法教授)。Voldemort 是一个优秀的分布式 KV 系统,但它有几个问题:

  1. 缺乏原生批处理集成:没有内置的 Full Push 机制,每次更新数据集都需要业务方自己实现数据加载流程
  2. 版本管理简陋:没有多版本并行(backup/current/future)的概念,升级数据集等于停服
  3. 混合读写支持弱:无法优雅地处理批处理全量写入 + 流处理增量写入的混合场景
  4. 多区域部署不成熟:没有 CRDT 支持,跨机房部署需要业务层做大量协调

到 2018 年,LinkedIn 约 500 个 Voldemort Read-Only 使用场景全部迁移到了 Venice,利用其 Full Push 能力作为无缝替代。这验证了 Venice 的设计方向是对的。

1.3 Venice 的设计哲学

Venice 的设计哲学用一句话概括:「把异步写入做到极致,同时让在线读取足够快」

这不是一个通用数据库,而是一个为衍生数据场景深度优化的专用系统。它做了几个违背「通用性」的设计决策:

  • 不支持强一致性写入(所有写入都是异步的)
  • 不支持 SQL(完全 KV 接口 + Read Compute DSL)
  • 数据来源是固定的(Hadoop/Spark 批处理 + Samza 流处理)

这些「不做」反而让它在目标场景里无出其右。


二、核心概念:Store、Version 与 Partition

2.1 Store:数据集即 store

在 Venice 中,一个 Store 对应一个数据集,类似于关系型数据库中的一张表,或者 KV 数据库中的一个桶(bucket)。每个 Store 有以下关键属性:

store_name: "user_recommendation_features"
partition_count: 1024          # 水平分区数
replication_factor: 3         # 每分区副本数
write_computation_enabled: true # 是否支持 Write Compute
hybrid_store_enabled: true      # 是否为混合存储
rewind_time_in_seconds: 86400  # 混合存储回溯窗口(1天)

2.2 多版本并行(Multi-Version Architecture)

这是 Venice 最核心的设计之一。一个 Store 可以同时运行最多三个版本

┌──────────────────────────────────────────────────────────────┐
│  STORE: user_features                                        │
│                                                              │
│  ┌──────────────┐  ┌──────────────┐  ┌──────────────┐     │
│  │   Version 10 │  │   Version 11 │  │  Version 12  │     │
│  │   (Backup)   │  │   (Current)  │  │   (Future)   │     │
│  └──────────────┘  └──────────────┘  └──────────────┘     │
│       READ ──────────────────────────────►                 │
│                   ▲                                          │
│                   │                                          │
│              SWAP POINT                                      │
│                                                              │
│  写入路径: 同时写入 backup + current + future                 │
│  读取路径: 只读 current(原子切换,无停服)                  │
└──────────────────────────────────────────────────────────────┘

为什么需要三个版本同时存在?

这个设计解决了一个分布式系统经典问题:如何在不影响在线读取的情况下更新一个数据集?

传统做法:

  1. 停服
  2. 删除旧数据
  3. 加载新数据
  4. 重新上线

Venice 做法(Full Push):

  1. 新数据版本(Version 12)作为「Future」在后台加载
  2. 加载期间,读请求继续服务 Version 11(Current)
  3. Version 12 加载完成后,原子切换读取指向 Version 12
  4. Version 11 降级为 Backup(可回滚)
  5. 旧版本自动按策略清理

这个「三版本并行 + 原子切换」的机制,让 Venice 实现了零停机时间的数据集更新。对业务方来说,数据集升级就像魔法一样——请求没断,数据已经换了。

2.3 分区与副本

数据水平分区(Partition),每个分区在多个节点上复制以保证高可用。分区数在创建 Store 时指定,LinkedIn 的典型配置是 1024 或 2048 个分区,足以支持大规模水平扩展。


三、写路径:四种写入模式与混合负载

3.1 2×3 写入矩阵

Venice 的写入 API 可以从两个维度理解:数据来源 × 写入模式,形成完整的 2×3 矩阵:

Hadoop / SparkSamza / Flink / Online Producer
全量数据集替换Full Push JobStream Reprocessing Job
向现有数据集插入行Incremental Push JobStreaming Writes
更新现有行的某些列Incremental Push JobStreaming Writes
(doing partial updates)(doing partial updates)

六种组合全部支持,这在工程实现上并不简单。

3.2 Full Push Job:批量全量替换

Full Push 是最常用的写入方式。数据从 Hadoop/Spark 抽出,整块写入 Venice:

// Venice Push Job 配置示例(伪代码)
VenicePushJobConfig config = VenicePushJobConfig.builder()
    .setInputDataPath("hdfs://warehouse/user_features/dt=2026-07-29")
    .setVeniceStoreName("user_recommendation_features")
    .setPushJobType(PushJobType.INCREMENTAL)  // 或 INCREMENTAL
    .setEnablePartialUpdates(false)            // Full Push
    .setCompressionStrategy(CompressionStrategy.ZSTD)
    .build();

VenicePushJob pushJob = new VenicePushJob(config);
pushJob.run();

Full Push 的执行流程:

1. Push Job 启动,从 HDFS 读取数据
2. 数据分片,每片由一个 Partition Worker 处理
3. 数据写入 Future 版本的后台存储(不阻断读取)
4. Future 版本加载完成后,进入 "SWAP_AFTER_PUSH" 阶段
5. 原子切换:Current → Backup,Future → Current
6. 旧 Backup 按策略清理(或保留 N 个历史版本)

整个过程对在线读取完全透明,没有任何锁或停服窗口。

3.3 Incremental Push:增量插入

与 Full Push 不同,Incremental Push 向现有所有版本(backup、current、future)插入新数据,而不是替换整个数据集:

// Incremental Push:插入特定行
VeniceIncrementalPushConfig config = VeniceIncrementalPushConfig.builder()
    .setInputDataPath("hdfs://delta/user_features_new_rows")
    .setIncrementalPushVersion("incr_20260729_01")
    .build();

典型使用场景:多源数据汇合——上游有多个独立的 ETL 任务,各自负责不同列的更新,通过 Incremental Push 汇入同一个 Store。

3.4 Streaming Writes:流式实时写入

通过 Apache Samza(或 Flink)进行流式写入:

// Samza 作业写入 Venice
SamzaVeniceWriter<Double> veniceWriter = new SamzaVeniceWriter<>(
    "user_behavior_stream",  // Store 名
    VeniceSystemFactory.INCREMENTAL_PUSH_SYSTEM_NAME
);

context.out().send(
    veniceWriter.write("user_12345",      // key
        userBehaviorUpdate,               // value(部分更新)
        "stream_realtime_v1"              // version name
    )
);

Streaming Writes 与 Incremental Push 一样,写入所有现有版本,保证数据一致性。

3.5 Write Compute:声明式部分更新(重点!)

这是 Venice 最强大的能力之一,也是它区别于大多数 KV 存储的关键功能。

在传统 KV 系统中,如果你想「只更新用户记录中的某个字段」,你需要这样做:

// 传统 KV 系统的做法(Read-Modify-Write,不可行)
UserRecord record = kvStore.get("user_12345"); // 读取
record.setLastLoginTime(now());               // 修改
kvStore.put("user_12345", record);             // 写回(两步操作)

这个流程在 Venice 中完全不可行,因为所有写入都是异步的。你无法在写入时做「读-改-写」。

Venice 的解法是:把修改逻辑下沉到服务器端,通过声明式的 Write Compute 操作表达增量

// Venice Write Compute:只更新特定字段
VeniceWriteCompute writeCompute = VeniceWriteCompute.build()
    .setKey("user_12345")
    // Partial Update:只更新 last_login_time 列
    .addPartialUpdate("last_login_time", currentTimestamp)
    // Collection Merge:向 tags 集合添加元素
    .addToCollection("tags", "premium_member")
    .addToCollection("viewed_items", "item_789")
    // 从集合中删除元素
    .removeFromCollection("abandoned_cart", "item_456");

veniceWriter.write(writeCompute);

服务器端执行这些操作时,不需要知道行的完整内容,只需要知道「哪个字段加什么值」或「哪个集合加/减什么元素」。这不仅解决了异步写入的问题,还带来了巨大的性能收益:

  • 网络传输减少:只发送变更,不发送整行
  • 服务器端原子执行:避免并发写入时的数据竞争
  • 支持高并发多写方:多个 Samza 作业可以同时向同一个 Store 写入,无需协调

Collection Merging 还支持 Map 结构的增删操作,对于特征存储场景特别有用——每个用户的特征向量可以作为一个 Map,训练 pipeline 可以独立更新不同维度的特征。

3.6 混合负载:Batch + Stream 的优雅融合

最复杂的场景来了:同时有全量批处理和实时流写入,如何保证数据一致性?

这就是 Hybrid Store 的用武之地。配置了 hybrid_store_enabled=true 的 Store 允许同时接收:

  • Full Push(全量替换)
  • Incremental Push / Streaming Writes(增量更新)

两者如何融合?Venice 使用 Rewind Time 机制

# 混合存储配置
store_hybrid_store_enabled: true
rewind_time_in_seconds: 86400  # 回溯1天内的实时写入

Hybrid 的执行流程:

时间线:
Day 7 14:00 ─── Full Push Version 12 开始后台加载
Day 7 14:30 ─── Full Push Version 12 加载完成
Day 7 14:30 ─── REPLAY 阶段启动
Day 7 14:30 ─── 从 Day 6 14:30 开始的实时流写入被回放(replay)到 Version 12
Day 7 15:00 ─── REPLAY 追上实时写入进度
Day 7 15:00 ─── 原子切换:Current → Version 12
Day 7 15:00 ─── 旧 Version 11 成为 Backup

核心洞察:Full Push 提供历史数据的完整快照,Streaming Writes 提供快照之后的增量变化。Rewind Time 定义了「增量」的窗口——超过这个窗口的实时数据,在 Full Push 时不会被回放,因为那部分数据已经被全量覆盖了。

这个设计与 Lambda 架构的思想一致,但 Venice 在工程层面做了更干净的实现(不需要维护两套独立系统)。

截至开源时,LinkedIn 已有 200+ 混合 Store 在生产环境运行


四、读路径:从 2 跳网络请求到 0 跳本地存储

4.1 三种客户端:性能与资源消耗的权衡

Venice 提供了四种客户端,分属三个性能档位:

客户端类型网络跳数P99 延迟状态典型场景
Thin Client2 跳< 10 ms无状态通用场景,灵活扩展
Fast Client1 跳< 2 ms轻量(路由元数据)低延迟,高吞吐
Da Vinci(RAM)0 跳< 10 μs全量内存超低延迟,资源密集
Da Vinci(SSD)0 跳< 1 ms全量 SSD低延迟,大数据集(内存放不下)

Thin Client(2 跳):

客户端 → Router(路由层)→ Storage Node(存储层)
         ↑
    Router 维护完整路由表,
    知道每个 Partition 在哪个 Storage Node

Thin Client 是完全无状态的,适合 Kubernetes 无状态部署。所有客户端共享同一套路由 API,成本/性能调优不需要改业务代码——换一个客户端类型即可。

Fast Client(1 跳):

客户端 → Storage Node(直连)
         ↑
    客户端内置路由缓存,
    知道 Partition → Storage Node 的映射

Fast Client 是分区感知的(partition-aware),通过缓存路由元数据省掉 Router 这一跳。将延迟从 < 10ms 压缩到 < 2ms。

Da Vinci Client(0 跳):

客户端 ──本地存储── Storage Node(同进程)

Da Vinci 是一个嵌入式存储引擎,直接把 Partition 数据加载到进程本地(内存或 SSD)。读取完全在本地完成,延迟压到微秒级。

这三种客户端共用同一套 Read API

// Venice Read API(统一接口)
VeniceClient client = new FastClient(storeName);  // 换成 DaVinciClient 即可

// Single Get
byte[] value = client.get("user_12345");

// Batch Get
Map<String, byte[]> values = client.batchGet(keys);

// Read Compute(服务器端计算)
ReadCompute readCompute = ReadCompute.newBuilder()
    .select("embedding_vector", "user_score")
    .compute(CosineSimilarity.of("embedding_vector", queryVector))
    .build();
byte[] result = client.get("user_12345", readCompute);

4.2 Read Compute:把计算推送到数据所在位置

这是 Venice 最具技术深度的读取能力。Read Compute 是一种声明式的数据处理 DSL,允许在服务器端执行以下操作:

// Read Compute 示例
ReadCompute compute = ReadCompute.newBuilder()
    // 字段投影:只返回需要的列
    .select("user_id", "embedding_vector", "last_updated")
    // 点积:向量内积
    .compute(DotProduct.of("embedding_vector", queryVector))
    // 余弦相似度
    .compute(CosineSimilarity.of("embedding_vector", queryVector))
    // Hadamard 积(逐元素乘法)
    .compute(HadamardProduct.of("embedding_vector", queryVector))
    // 集合计数
    .compute(CollectionCount.of("user_tags"))
    .build();

为什么这个能力如此重要?

以 LinkedIn 的 People You May Know(PYMK) 为例:这个推荐系统需要根据用户 embedding 向量,从数十亿用户中找出最相似的 Top-K 候选。

没有 Read Compute 的做法:

1. 从 Venice 读取所有候选用户的 embedding(亿级向量,每个 512 维 float)
2. 本地计算余弦相似度(数据传输量:512 维 × 亿级行 = 天文数字)
3. 返回 Top-K

问题:网络传输量巨大,P99 延迟爆炸

使用 Read Compute 的做法:

// 只请求 Top-10 最相似的用户 ID 和相似度分数
ReadCompute compute = ReadCompute.newBuilder()
    .select("user_id")
    .compute(TopKCosineSimilarity.of(
        "embedding_vector",    // 存储的向量
        myEmbedding,           // 查询向量
        10                      // Top-K
    ))
    .build();
Result result = client.get("user_12345", compute);
// 服务器端计算,只返回 10 个结果

这相当于把推荐系统的向量检索逻辑下推到 Venice 存储层,极大地减少了网络传输。LinkedIn 从 2019 年开始将这个能力用于 PYMK 的在线深度学习推理,实现了水平扩展的在线向量相似度计算——这是一个在工业界相当罕见的成就。


五、多区域多活:CRDT 如何解决跨机房冲突

5.1 为什么需要 Active-Active

LinkedIn 是全球化的服务,必须在多个地理区域部署。用户请求通常路由到最近的区域,但如果区域 A 的用户需要读取区域 B 写入的数据怎么办?

传统方案:主从复制(Master-Slave)

  • 优点:实现简单,没有冲突
  • 缺点:跨区域读取延迟高(主从同步有延迟)

Venice 的选择:Active-Active 多区域复制

  • 每个区域都可以写入
  • 读取就近访问本地副本
  • 跨区域写入可能产生冲突

5.2 CRDT-based 冲突解决

既然多个区域同时写入同一份数据,冲突不可避免。Venice 使用 CRDT(Conflict-free Replicated Data Types) 来解决。

CRDT 的核心思想:设计数据结构,使得任意顺序的合并操作都能得到确定的结果

Venice 的 Write Compute 操作天然支持 CRDT:

  • Last-Write-Wins(LWW)Register:标量字段(字符串、数值),按时间戳决定胜者
  • Grow-Only Set(G-Set):集合只增不减(添加标签、行为记录)
  • Add-Wins Set(AW-Set):集合支持添加和删除,添加操作优先

以用户标签集合为例:

区域 A(UTC+0):addToSet("user_123", "premium")
区域 B(UTC-5):addToSet("user_123", "verified")

CRDT 合并结果:{"premium", "verified"}(两个标签都保留)

这种合并方式是确定性的,无论网络延迟如何、无论合并顺序如何,最终状态一致。这比传统的主从复制简单得多,因为不需要「谁先生效」的业务层决策。

5.3 多区域部署架构

┌──────────────────────────────────────────────────────────────┐
│                    Global View                                │
│                                                              │
│  Region US-East          Region EU-Central      Region APAC │
│  ┌──────────────┐        ┌──────────────┐      ┌──────────────┐
│  │ Venice       │◄──────►│ Venice       │◄────►│ Venice       │
│  │ Cluster      │  CRDT  │ Cluster      │  CRDT│ Cluster      │
│  │              │ Sync   │              │ Sync │              │
│  │ Write locally │────────│ Write locally │──────│ Write locally │
│  │ Read locally  │        │ Read locally  │      │ Read locally  │
│  └──────────────┘        └──────────────┘      └──────────────┘
│                                                              │
│  Samza + Kafka(区域独立数据源)                            │
└──────────────────────────────────────────────────────────────┘

每个区域维护独立的 Samza + Kafka 数据源,通过 CRDT 机制跨区域同步。


六、Venice 与 Feature Store:ML 推理的最佳拍档

6.1 Feature Store 的存储困境

现代 ML 系统通常分为离线训练在线推理两个阶段:

  • 离线训练:用 Spark/Hadoop 跑批处理,生成特征数据,训练模型
  • 在线推理:模型上线,实时查询特征,进行预测

问题是:训练阶段用到的特征,必须和推理阶段用到的特征完全一致。特征定义变了,模型就废了。

Feature Store 的职责就是统一管理特征的定义、版本和 Serving,确保训练和推理使用同一套特征。

6.2 Venice 作为 Feature Store 的在线存储层

LinkedIn 的 Feature Store Feathr 使用 Venice 作为在线存储层:

┌─────────────────────────────────────────────────────────────────┐
│                     Feathr Feature Store                         │
│                                                                 │
│  ┌─────────────┐    ┌─────────────┐    ┌─────────────────────┐ │
│  │  离线训练    │    │  特征注册表  │    │   在线推理          │ │
│  │  Spark Job  │───►│   Registry  │───►│   Model Serving     │ │
│  └─────────────┘    └─────────────┘    └─────────────────────┘ │
│         │                                      │                 │
│         ▼                                      ▼                 │
│  ┌─────────────────────────────────────────────────────────┐    │
│  │           Venice(Feature Store 在线存储层)              │    │
│  │                                                         │    │
│  │  Batch Push ← Hadoop/Spark 训练特征                     │    │
│  │  Streaming Writes ← Samza 实时特征更新                   │    │
│  │                                                         │    │
│  │  Read Compute ← 在线推理查询(向量相似度)                │    │
│  └─────────────────────────────────────────────────────────┘    │
└─────────────────────────────────────────────────────────────────┘

训练时,Spark 作业将特征写入 HDFS,通过 Full Push 导入 Venice。
推理时,Model Serving 通过 Da Vinci Client(0 跳,< 10 μs)查询特征。

6.3 为什么 Venice 比 Redis 更适合 Feature Store?

有人可能会问:为什么不直接用 Redis?

维度RedisVenice
数据来源需要手动同步原生集成 Hadoop/Spark/Samza
全量更新需停服或双写Full Push 原子切换,零停服
批量写入吞吐低(主从同步)峰值 50 GB/s
特征版本管理需业务层实现内置多版本并行
向量相似度需外部插件原生 Read Compute(点积/余弦)
多区域复制Redis Cluster原生 CRDT Active-Active
数据量受内存限制Da Vinci 支持 SSD 扩展

Venice 的每一个设计决策,都在解决 Feature Store 场景的真实痛点。


七、生产规模与性能调优

7.1 LinkedIn 生产数据(2022 年开源时)

指标数值
数据集数量1800+
应用数量300+
每日写入量1.2 PB
日写入行数3 万亿+
平均写入吞吐14 GB/s, 39M 行/s
峰值写入吞吐50 GB/s, 113M 行/s
单服务器写入吞吐(节流配置)200 MB/s, 40万行/s
单服务器 Serving QPS~20 万/秒
峰值读取 QPS(Batch Get)4500 万 Key 查找/s
Read Compute QPS(峰值)4600 万 Key 查找/s

7.2 Throttling 策略

Venice 在写入路径上实现了精细的 Throttling,以防止写入过载影响在线读取:

// 写入节流配置(per-server)
VeniceWriterConfig config = VeniceWriterConfig.builder()
    .setMaxWriteThroughputBytesPerSecond(200 * 1024 * 1024)  // 200 MB/s
    .setMaxWriteThroughputRowsPerSecond(400_000)              // 40 万行/s
    .build();

即便在 200 MB/s + 40 万行/s 的写入吞吐下,服务器仍能正常 Serving 在线读取请求,没有抖动。这是通过写入队列优先级控制和后台 I/O 调度实现的。

7.3 数据淘汰与版本管理

# 版本淘汰策略
number_of_versions_to_keep: 3         # 保留最近3个版本
retention_time_in_seconds: 604800     # 或保留7天(以先到者为准)

Venice 自动管理版本生命周期:超过保留策略的旧版本自动删除,释放存储空间。业务方不需要手动清理。


八、快速上手:写一个 Venice Feature Store 客户端

8.1 添加依赖(Maven)

<repositories>
    <repository>
        <id>venice-jfrog</id>
        <name>VeniceJFrog</name>
        <url>https://linkedin.jfrog.io/artifactory/venice</url>
    </repository>
</repositories>

<dependencies>
    <dependency>
        <groupId>com.linkedin.venice</groupId>
        <artifactId>venice-client</artifactId>
        <version>0.4.455</version>
    </dependency>
</dependencies>

8.2 初始化客户端(Fast Client)

import com.linkedin.venite.client.VeniceClient;
import com.linkedin.venite.client.VeniceClientFactory;

public class VeniceFeatureStoreDemo {
    public static void main(String[] args) {
        // 创建 Fast Client(1跳,低延迟)
        VeniceClientFactory factory = new VeniceClientFactory.Builder()
            .setVeniceServers("venice-cluster.example.com:8080")
            .setRouterUrls("router1.example.com:8080,router2.example.com:8080")
            .build();

        VeniceClient client = factory.getVeniceClient("user_features");

        // ========== 写入操作(通过 Push Job,这里演示 Write Compute)==========
        VeniceWriter<String, byte[]> writer = client.getWriter();
        
        // Partial Update:只更新 embedding 向量
        writer.writePartialUpdate(
            "user_10001",
            Map.of(
                "embedding_vector", serializeVector(new float[]{0.1f, 0.3f, 0.5f}),
                "last_updated", System.currentTimeMillis()
            ),
            "stream_v1"
        );

        // ========== 读取操作 ===========
        // Single Get
        byte[] value = client.get("user_10001");
        UserFeature feature = deserialize(value);
        System.out.println("User: " + feature.userId);

        // Batch Get(批量拉取多个用户的特征)
        List<String> userIds = Arrays.asList("user_10001", "user_10002", "user_10003");
        Map<String, byte[]> batchResult = client.batchGet(userIds);

        // Read Compute(服务器端向量相似度计算)
        float[] myEmbedding = new float[]{0.2f, 0.4f, 0.6f};
        ReadCompute compute = ReadCompute.newBuilder()
            .select("user_id")
            .compute(TopKCosineSimilarity.of("embedding_vector", myEmbedding, 10))
            .build();

        ReadComputeResult result = client.get("user_10001", compute);
        List<ScoredUser> topSimilar = result.getTopKSimilarUsers();
        
        for (ScoredUser user : topSimilar) {
            System.out.printf("User: %s, Score: %.4f%n",
                user.userId, user.similarityScore);
        }
    }
}

8.3 数据类型支持

Venice 的 Value 支持以下数据类型:

// Venice 支持的 Value Schema 类型
VeniceSchema schema = VeniceSchema.builder()
    .setPrimaryKey("user_id")
    .addField("user_id", STRING)           // 主键
    .addField("embedding", FLOAT_ARRAY)    // float[](向量)
    .addField("tags", STRING_SET)           // 集合类型(Write Compute 支持)
    .addField("metadata", MAP)              // Map<String, String>
    .addField("score", DOUBLE)              // 数值
    .addField("is_active", BOOLEAN)         // 布尔
    .build();

集合类型(SET、MAP)是 Write Compute partial update 的基础,支持 addToSetremoveFromSetaddToMap 等原子操作。


九、局限性:Venice 不是银弹

作为一个严肃的技术文章,坦诚地说明 Venice 的局限性是必要的:

9.1 不适合的场景

  1. 强一致性写入需求:如果你的业务需要用户在点击「保存」后立即看到更新,Venice 不适合(写入是异步的,延迟从秒级到分钟级不等)
  2. 复杂查询:Venice 只有 KV + Read Compute,没有 SQL,不支持范围查询、聚合查询、JOIN
  3. 小规模数据:Venice 的运维复杂度较高(多组件、多版本管理),对于简单场景是杀鸡用牛刀
  4. 非衍生数据:主数据(用户主档、交易记录)应该用 OLTP 数据库,Venice 只接受衍生数据

9.2 开源生态的挑战

截至目前(2026年),Venice 的 Maven 依赖尚未发布到 Maven Central,需要额外配置 JFrog 仓库。虽然这只是个小麻烦,但确实增加了试用门槛。


十、总结:Venice 教给我们什么

10.1 工程哲学

Venice 的成功源于一个朴素的工程原则:不要做一个什么都做的系统,做一个特定场景下无可替代的系统

  • 它放弃了强一致性,换来了 50 GB/s 的写入吞吐
  • 它放弃了 SQL,换来了毫秒级的向量 Read Compute
  • 它放弃了通用性,换来了 Feature Store 场景的端到端优化

这是分布式系统设计中最难的一课:知道不做什么,比知道做什么更重要

10.2 衍生数据平台的趋势

随着大模型和 AI Native 应用的爆发,Feature Store 和衍生数据平台的价值正在被重新定义。传统的「数据库 + ETL + 应用」架构,正在被「原生衍生数据平台 + 在线推理」架构取代。Venice 作为这个领域的先驱实践,其架构设计值得所有做 AI Infra 的工程师深入研究。

10.3 给工程师的建议

如果你在考虑引入 Venice 或类似的衍生数据平台,以下问题值得认真思考:

  1. 你的数据有多少比例是「衍生数据」(由计算生成,而非用户直接产生)?
  2. 你的离线训练和在线推理是否使用同一套特征?一致性如何保证?
  3. 你的向量相似度计算目前在哪里做?是否有性能瓶颈?
  4. 你的数据更新频率如何?是否适合异步写入模型?

如果这些问题中,你对多个问题的答案是「痛点明显」,那么 Venice(或 Feathr + Venice)的组合值得深入评估。


参考资源:

推荐文章

HTML + CSS 实现微信钱包界面
2024-11-18 14:59:25 +0800 CST
Vue3中哪些API被废弃了?
2024-11-17 04:17:22 +0800 CST
在 Docker 中部署 Vue 开发环境
2024-11-18 15:04:41 +0800 CST
CSS 特效与资源推荐
2024-11-19 00:43:31 +0800 CST
Nginx 防盗链配置
2024-11-19 07:52:58 +0800 CST
程序员茄子在线接单