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 系统,但它有几个问题:
- 缺乏原生批处理集成:没有内置的 Full Push 机制,每次更新数据集都需要业务方自己实现数据加载流程
- 版本管理简陋:没有多版本并行(backup/current/future)的概念,升级数据集等于停服
- 混合读写支持弱:无法优雅地处理批处理全量写入 + 流处理增量写入的混合场景
- 多区域部署不成熟:没有 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(原子切换,无停服) │
└──────────────────────────────────────────────────────────────┘
为什么需要三个版本同时存在?
这个设计解决了一个分布式系统经典问题:如何在不影响在线读取的情况下更新一个数据集?
传统做法:
- 停服
- 删除旧数据
- 加载新数据
- 重新上线
Venice 做法(Full Push):
- 新数据版本(Version 12)作为「Future」在后台加载
- 加载期间,读请求继续服务 Version 11(Current)
- Version 12 加载完成后,原子切换读取指向 Version 12
- Version 11 降级为 Backup(可回滚)
- 旧版本自动按策略清理
这个「三版本并行 + 原子切换」的机制,让 Venice 实现了零停机时间的数据集更新。对业务方来说,数据集升级就像魔法一样——请求没断,数据已经换了。
2.3 分区与副本
数据水平分区(Partition),每个分区在多个节点上复制以保证高可用。分区数在创建 Store 时指定,LinkedIn 的典型配置是 1024 或 2048 个分区,足以支持大规模水平扩展。
三、写路径:四种写入模式与混合负载
3.1 2×3 写入矩阵
Venice 的写入 API 可以从两个维度理解:数据来源 × 写入模式,形成完整的 2×3 矩阵:
| Hadoop / Spark | Samza / Flink / Online Producer | |
|---|---|---|
| 全量数据集替换 | Full Push Job | Stream Reprocessing Job |
| 向现有数据集插入行 | Incremental Push Job | Streaming Writes |
| 更新现有行的某些列 | Incremental Push Job | Streaming 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 Client | 2 跳 | < 10 ms | 无状态 | 通用场景,灵活扩展 |
| Fast Client | 1 跳 | < 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?
| 维度 | Redis | Venice |
|---|---|---|
| 数据来源 | 需要手动同步 | 原生集成 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 的基础,支持 addToSet、removeFromSet、addToMap 等原子操作。
九、局限性:Venice 不是银弹
作为一个严肃的技术文章,坦诚地说明 Venice 的局限性是必要的:
9.1 不适合的场景
- 强一致性写入需求:如果你的业务需要用户在点击「保存」后立即看到更新,Venice 不适合(写入是异步的,延迟从秒级到分钟级不等)
- 复杂查询:Venice 只有 KV + Read Compute,没有 SQL,不支持范围查询、聚合查询、JOIN
- 小规模数据:Venice 的运维复杂度较高(多组件、多版本管理),对于简单场景是杀鸡用牛刀
- 非衍生数据:主数据(用户主档、交易记录)应该用 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 或类似的衍生数据平台,以下问题值得认真思考:
- 你的数据有多少比例是「衍生数据」(由计算生成,而非用户直接产生)?
- 你的离线训练和在线推理是否使用同一套特征?一致性如何保证?
- 你的向量相似度计算目前在哪里做?是否有性能瓶颈?
- 你的数据更新频率如何?是否适合异步写入模型?
如果这些问题中,你对多个问题的答案是「痛点明显」,那么 Venice(或 Feathr + Venice)的组合值得深入评估。
参考资源:
- GitHub: linkedin/venice
- 官网文档: venicedb.org
- LinkedIn 工程博客: Open Sourcing Venice
- Conference Talk: QCon AI - Scaling Deep Learning in Production