Pixels-Retina: 面向数据湖的实时变更同步框架深度解读
Pixels 是来自中国人民大学卞老师团队的开源项目: https://github.com/pixelsdb/pixels
今天聊一聊 Pixels 项目中的 Pixels-Retina 组件.
一、项目概述与定位
1.1 什么是 Pixels-Retina
Pixels-Retina 是 Pixels 项目生态中的实时数据同步框架,专门用于将变更数据捕获(CDC)源的数据变更操作,以镜像事务(Mirror Transaction) 的形式重放到列式表数据上。
用一句话概括: Retina 实现了在数据湖上达到数据库级别的实时变更能力,同时保留了分析型查询的大规模数据读取能力。
1.2 解决的核心痛点
传统数据湖(如 Iceberg、Paimon)采用 Merge-On-Read (MoR) 机制,存在以下问题:
| 数据新鲜度差 | |
| 变更处理吞吐低 | |
| 查询性能不稳定 | |
| 资源成本高 |
Retina 的核心价值: 10毫秒级数据新鲜度 + 320万行/秒的变更重放吞吐量
1.3 适用场景
┌─────────────────────────────────────────────────────────────────┐
│ 适用场景 │
├─────────────────────────────────────────────────────────────────┤
│ • 实时数据仓库:需要近实时同步的上游 CDC 数据 │
│ • 湖仓一体:既要支持高频写入,又要支持大规模分析查询 │
│ • CDC 复制:从 MySQL/PostgreSQL 等数据库实时同步到湖里 │
│ • 审计日志:记录所有数据变更,支持时间点回溯 │
│ • 变更数据流:需要追踪数据全量变更历史的业务 │
└─────────────────────────────────────────────────────────────────┘
二、架构设计总览
2.1 Pixels 整体架构
根据 DeepWiki 提供的架构信息,Pixels 系统分为以下层次:
┌─────────────────────────────────────────────────────────────────┐
│ 客户端层 │
│ ┌─────────────┐ ┌─────────────┐ ┌─────────────┐ │
│ │ Trino │ │ DuckDB │ │ pixels-cli │ │
│ └─────────────┘ └─────────────┘ └─────────────┘ │
└────────────────────────────┬────────────────────────────────────┘
│
┌────────────────────────────▼────────────────────────────────────┐
│ 服务层 (pixels-daemon) │
│ ┌──────────┐ ┌──────────┐ ┌──────────┐ ┌──────────┐ │
│ │Metadata │ │ Trans │ │ Retina │ │ Index │ │
│ │ Server │ │ Server │ │ Server │ │ Server │ │
│ └──────────┘ └──────────┘ └──────────┘ └──────────┘ │
└────────────────────────────┬────────────────────────────────────┘
│
┌────────────────────────────▼────────────────────────────────────┐
│ 核心库 (pixels-*) │
│ ┌──────────┐ ┌──────────┐ ┌──────────┐ │
│ │pixels- │ │pixels- │ │pixels- │ │
│ │core │ │cache │ │retina │ │
│ └──────────┘ └──────────┘ └──────────┘ │
└────────────────────────────┬────────────────────────────────────┘
│
┌───────────────────┼───────────────────┐
│ │ │
▼ ▼ ▼
┌─────────────────┐ ┌─────────────────┐ ┌─────────────────┐
│ 状态管理层 │ │ 存储后端 │ │ 索引后端 │
│ etcd + MySQL │ │ S3/HDFS/Redis │ │ RocksDB/SQLite │
└─────────────────┘ └─────────────────┘ └─────────────────┘
2.2 Retina 在架构中的位置
┌──────────────────────────────────────────────────────────────────┐
│ CDC 数据流 │
│ │
│ ┌─────────┐ ┌─────────┐ ┌─────────────────┐ │
│ │ MySQL │───▶│ Debezium│───▶│ Kafka │ │
│ │ (源库) │ │ (CDC) │ │ (消息队列) │ │
│ └─────────┘ └─────────┘ └────────┬────────┘ │
│ │ │
│ ▼ │
│ ┌─────────────────────┐ │
│ │ pixels-sink │ │
│ │ (独立服务,消费Kafka) │ │
│ └──────────┬──────────┘ │
│ │ gRPC RPC │
│ ▼ │
│ ┌─────────────────────┐ │
│ │ pixels-retina │ │
│ │ (镜像事务重放) │ │
│ └──────────┬──────────┘ │
│ │ │
│ ▼ │
│ ┌─────────────────────────────────────┐ │
│ │ 列式文件存储 │ │
│ │ ┌─────────┐ ┌─────────┐ │ │
│ │ │ ordered │ │ compact │ │ │
│ │ │ 路径 │ │ 路径 │ │ │
│ │ └─────────┘ └─────────┘ │ │
│ └─────────────────────────────────────┘ │
└──────────────────────────────────────────────────────────────────┘
注意: pixels-sink 是独立项目, 不是 Flink connector
2.3 Retina 的核心组件
pixels-retina/
├── Java 层
│ ├── RetinaResourceManager # 核心资源管理器(单例)
│ │ ├── rgVisibilityMap # 行组可见性映射
│ │ ├── pixelsWriteBufferMap # 写缓冲区映射
│ │ └── GC 执行器 # 内存 GC + 存储 GC
│ │
│ ├── StorageGarbageCollector # 存储垃圾回收器
│ │ ├── FileCandidate # GC 候选文件
│ │ ├── FileGroup # 文件分组
│ │ └── RewriteResult # 重写结果
│ │
│ ├── RGVisibility # 行组可见性(JNI 调用 C++)
│ ├── PixelsWriteBuffer # 写缓冲区
│ │ ├── MemTable # 内存表
│ │ ├── ObjectEntry # 对象条目
│ │ └── SuperVersion # 当前版本视图
│ │
│ └── ObjectStorageManager # 对象存储管理
│
└── C++ 层 (cpp/pixels-retina/)
├── TileVisibility # 行组可见性核心实现
├── VersionedData # 版本化数据(COW)
└── DeleteIndexBlock # 删除链
三、底层原理深度剖析
3.1 MVCC 机制:轻量级多版本并发控制
Retina 采用了轻量级 MVCC 而非传统数据库的复杂事务机制:
传统数据库 MVCC vs Retina MVCC
┌─────────────────────────────────────────────────────────────────┐
│ 传统数据库 MVCC │
│ • 完整的 Undo Log / Redo Log │
│ • 复杂的锁机制(行锁、表锁) │
│ • 事务生命周期管理 │
│ • 开销大,不适合高并发场景 │
└─────────────────────────────────────────────────────────────────┘
┌─────────────────────────────────────────────────────────────────┐
│ Retina 轻量级 MVCC │
│ • 无 Undo Log,仅保留删除链 │
│ • Copy-On-Write 版本数据 │
│ • 基于时间戳的可见性判断 │
│ • 无锁读,仅在写入时轻量同步 │
└─────────────────────────────────────────────────────────────────┘
核心数据结构:
VersionedData:存储基础位图和基准时间戳
structVersionedData {
long* baseBitmap; // 基础可见性位图
long baseTimestamp; // 基准时间戳
};DeleteIndexBlock:删除链节点
structDeleteIndexBlock {
uint64_t* entries; // 打包的 (rowId, timestamp) 对
DeleteIndexBlock* next;
};TileVisibility:行组可见性管理器
classTileVisibility {
VersionedData* currentVersion; // 当前版本(COW)
DeleteIndexBlock* deleteChain; // 删除链
std::atomic<int> refCount; // 引用计数
};
3.2 可见性位图(Visibility Bitmap)
位图结构:每个行组维护一个位图,位表示该行是否被删除
位图示意 (8行举例):
Row: [0] [1] [2] [3] [4] [5] [6] [7]
Bitmap: 0 0 1 0 0 1 0 0
▲ ▲
已删除 已删除
含义:
• bit = 0:行可见(未被删除)
• bit = 1:行已删除
获取指定时间戳的可见性位图:
getVisibilityBitmap(timestamp):
1. 复制 baseBitmap
2. 遍历删除链:
对每个 (rowId, ts) 对:
如果 ts <= timestamp,则在位图中标记该行为已删除
3. 返回处理后的位图
3.3 删除链(Deletion Chain)
删除链是解决 COW 写放大 问题的关键设计:
写入场景时序:
时间 T1: DELETE row_5
├── 创建 DeleteIndexBlock: [(row_5, T1)]
└── baseBitmap 不变
时间 T2: DELETE row_3
├── 创建新 DeleteIndexBlock: [(row_3, T2)]
└── 链接到 T1 的 block
时间 T3: GET_VISIBILITY(ts=T2.5)
├── 复制 baseBitmap
├── 应用 T1 的删除: row_5 被标记
└── T2 的删除不应用(因为 T2 > T2.5)
时间 T4: GC(ts=100)
├── 将 T1 的删除合并到 baseBitmap
├── 释放 T1 的 DeleteIndexBlock
└── T2 的 block 成为链首
3.4 写入缓冲:MemTable + SuperVersion
数据写入路径:
CDC 数据
│
▼
┌─────────────────────────────────────────────┐
│ PixelsWriteBuffer │
│ │
│ ┌─────────────┐ │
│ │ active │ ◄── 当前写入的 MemTable │
│ │ MemTable │ │
│ └──────┬──────┘ │
│ │ 满了就切换 │
│ ┌──────▼──────┐ │
│ │ immutable │ ◄── 等待刷盘的 MemTable │
│ │ MemTables │ (可能多个) │
│ └──────┬──────┘ │
│ │ 异步刷盘 │
│ ┌──────▼──────┐ │
│ │ ObjectEntry │ ◄── 已刷到对象存储的数据 │
│ │ (引用计数) │ │
│ └─────────────┘ │
└─────────────────────────────────────────────┘
│
│ 定期合并
▼
┌─────────────┐
│ Pixels File │ (列式存储文件)
└─────────────┘
SuperVersion:提供数据的一致性视图
classSuperVersion{
MemTable activeMemTable; // 当前活跃的内存表
List<MemTable> immutableMemTables; // 不可变内存表(正在刷盘)
List<ObjectEntry> objectEntries; // 已刷到对象存储的条目
long versionId; // 版本号
}
3.5 GC 机制:内存 GC + 存储 GC
Retina 的 GC 分为两个层次:
3.5.1 内存 GC(Memory GC)
目标:压缩删除链,减少内存占用
Memory GC 执行流程 (RetinaResourceManager.runGC):
1. 获取安全 GC 时间戳 (TransService.getSafeGcTimestamp)
2. 遍历 rgVisibilityMap:
对每个 RGVisibility:
a. 调用 garbageCollect(timestamp)
- 将 timestamp 之前的删除合并到 baseBitmap
- 释放不再需要的 DeleteIndexBlock
b. 收集文件级统计 (totalRows, invalidCount)
3. 创建检查点(阻塞式,确保数据不丢失)
4. 传递给 Storage GC 候选文件列表
3.5.2 存储 GC(Storage GC)
目标:重写高删除率的存储文件,回收物理空间
Storage GC 核心逻辑 (StorageGarbageCollector):
输入:候选文件列表(删除率 > 阈值)
1. 文件分组:
按 (tableId, virtualNodeId) 分组
按删除率降序排序
2. 对每个文件组:
a. 读取旧文件,应用可见性过滤
b. 写入新文件(REGULAR 状态)
c. 注册双向写(dual-write):
- 新文件删除 → 同步到旧文件
- 旧文件删除 → 同步到新文件
d. 同步可见性链
e. 同步索引
f. 原子切换(TEMPORARY → REGULAR)
g. 延迟删除旧文件(grace period)
四、与竞品对比
4.1 技术对比表
| Pixels-Retina | Apache Iceberg | Apache Paimon | |
|---|---|---|---|
| 数据新鲜度 | |||
| 变更粒度 | |||
| MVCC 机制 | |||
| 合并策略 | |||
| 索引支持 | |||
| 存储格式 | |||
| CDC 支持 | |||
| 查询引擎 |
4.2 架构差异详解
Iceberg MoR 流程:
┌─────────────────────────────────────────────────────────────────┐
│ Write Path: │
│ INSERT → 写入新 Data File │
│ DELETE → 写入 Delete File (position delete / equality delete) │
│ │
│ Read Path: │
│ 读取 Data File + 应用 Delete File → Merge → Result │
│ 问题: 每次查询都需要 Merge,开销大 │
└─────────────────────────────────────────────────────────────────┘
Pixels-Retina 流程:
┌─────────────────────────────────────────────────────────────────┐
│ Write Path: │
│ CDC → PixelsWriteBuffer → 刷盘 → Pixels File │
│ 删除 → 更新可见性位图(内存操作) │
│ │
│ Read Path: │
│ 直接读取 Pixels File + 可见性位图过滤 → Result │
│ 优势: 无需 Merge,读取即可见 │
└─────────────────────────────────────────────────────────────────┘
4.3 性能数据对比
根据项目文档和论文:
| 数据新鲜度 | |||
| 变更重放吞吐 | |||
| 查询延迟 | |||
| 存储成本 |
五、使用前后效果变化
5.1 使用前(传统 MoR 方案)
业务场景: 电商订单实时同步
┌─────────────────────────────────────────────────────────────────┐
│ MySQL (源) ──CDC──▶ Kafka ──▶ Flink ──▶ Iceberg/Paimon │
│ │ │
│ ▼ │
│ Merge-On-Read │
│ (读取时合并) │
│ │ │
│ ▼ │
│ 查询延迟: 1-5 分钟 │
│ (需等待批量合并) │
└─────────────────────────────────────────────────────────────────┘
问题:
• 订单状态更新后,分析报表看到的是几分钟前的数据
• 数据分析师抱怨数据"不准"
• 高峰期合并操作积压,导致延迟进一步恶化
5.2 使用后(Retina 方案)
业务场景: 电商订单实时同步
┌─────────────────────────────────────────────────────────────────┐
│ MySQL (源) │
│ │ │
│ ▼ │
│ Debezium (CDC) │
│ │ │
│ ▼ │
│ Kafka │
│ │ │
│ ▼ │
│ pixels-sink (独立服务) │
│ │ │
│ ▼ │
│ Pixels-Retina │
│ │ │
│ ▼ │
│ PixelsWriteBuffer ──可见性位图更新──▶ 实时可见 │
│ │
│ 查询延迟: 10ms │
└─────────────────────────────────────────────────────────────────┘
效果:
• 订单状态更新后,毫秒级可查
• 报表数据实时准确
• 无需等待批量合并,无积压
5.3 具体案例
案例: 实时销售大屏
使用前:
• 数据延迟: 3-5 分钟 (Debezium → Kafka → Flink → Iceberg/Paimon → Merge)
• 销售总额展示滞后,导致决策延迟
• 客服无法实时查询最新订单状态
使用后:
• 数据延迟: <100ms (Debezium → Kafka → pixels-sink → Retina)
• 销售大屏实时更新,误差 <1 秒
• 客服可即时查询最新订单
六、最佳使用实践
6.1 部署架构建议
生产环境推荐部署:
┌──────────────────────────────────────────────────────────────────┐
│ 推荐架构 │
│ │
│ ┌─────────────┐ │
│ │ etcd 集群 │ (3节点) - 分布式协调 │
│ └─────────────┘ │
│ │ │
│ ┌──────┴──────┐ │
│ │ MySQL 集群 │ - 元数据存储 │
│ └─────────────┘ │
│ │
│ ┌─────────────┐ ┌─────────────┐ ┌─────────────┐ │
│ │Retina Node 1│ │Retina Node 2│ │Retina Node N│ │
│ │(写入节点) │ │(写入节点) │ │(写入节点) │ │
│ └─────────────┘ └─────────────┘ └─────────────┘ │
│ │ │ │ │
│ └─────────────────┴─────────────────┘ │
│ │ │
│ ┌──────┴──────┐ │
│ │ S3 / HDFS │ - 数据存储 │
│ └─────────────┘ │
└──────────────────────────────────────────────────────────────────┘
6.2 配置优化
核心配置项 (pixels.properties) :
# Retina 配置
retina.enable=true
retina.checkpoint.dir=/data/pixels/checkpoint
retina.checkpoint.threads=4
# GC 配置
retina.gc.interval=60 # GC 间隔(秒)
retina.storage.gc.enabled=true
retina.storage.gc.threshold=0.2 # 删除率阈值 20%
retina.storage.gc.target.file.size=128MB # 目标文件大小
retina.storage.gc.max.files.per.group=10
retina.storage.gc.max.file.groups.per.run=5
retina.storage.gc.file.retire.delay.hours=24 # 延迟删除时间
# 写入缓冲配置
retina.buffer.memTable.size=65536 # MemTable 大小 (64K rows)
retina.buffer.flush.count=4 # 刷盘阈值
retina.buffer.flush.interval=10 # 刷盘间隔(秒)
# 虚拟节点数(建议与 CPU 核数一致)
node.virtual.num=16
6.3 CDC 接入最佳实践
重要澄清:Pixels 的 CDC 接入架构是 pixels-sink(独立服务)消费 Debezium 数据,然后通过 gRPC 发送给 Retina。不是直接用 Flink 写入 Retina。
完整数据流架构
MySQL → Debezium → Kafka → pixels-sink → gRPC RPC → pixels-retina
Debezium + Kafka 配置
# Debezium MySQL Source 配置
connector.class:io.debezium.connector.mysql.MySqlConnector
database.hostname:mysql-master
database.port:3306
database.user:cdc_user
database.password:***
database.server.id:184054
database.server.name:pixels-cdc
table.include.list:orders,customers,products
snapshot.mode:schema_only
# 关键配置:将变更发送到 Kafka
include.schema.changes:false
transforms:unwrap
transforms.unwrap.type:io.debezium.transforms.ExtractNewRecordState
# Kafka topic 配置
topic.prefix:pixels-cdc
pixels-sink 配置
pixels-sink 是独立部署的服务,配置示例:
# pixels-sink.properties
retina.server.address=retina-node-1:18890
retina.server.port=18890
kafka.bootstrap.servers=kafka-broker:9092
kafka.group.id=pixels-sink-consumer
kafka.topic.pattern=pixels-cdc\\..*
debezium.format=json
sink.parallelism=16
部署说明
pixels-sink 是一个独立的 Java 进程(不是 Flink Job) 它订阅 Kafka 中 Debezium 的 topic 消费 CDC 消息,重建为镜像事务 通过 gRPC 调用 pixels-retina 的 UpdateRecord接口
⚠️ 注意:
pixels-sink是独立项目 (https://github.com/pixelsdb/pixels-sink),需要单独部署和配置。
6.4 监控与告警
关键监控指标:
# Prometheus 告警规则
groups:
-name:retina-alerts
rules:
-alert:RetinaGCBacklog
expr:retina_gc_pending_items>10000
for:5m
annotations:
summary:"Retina GC 积压严重"
-alert:RetinaCheckpointLag
expr:time()-retina_checkpoint_timestamp>300
annotations:
summary:"检查点创建延迟"
-alert:RetinaWriteBufferFull
expr:retina_write_buffer_capacity/retina_write_buffer_used<0.2
annotations:
summary:"写入缓冲区即将满"
6.5 常见问题处理
七、关键技术问题解答
基于 DeepWiki 架构分析,关于 Retina 的关键问题:
Q1: Retina 如何保证读写并发?
A: Retina 通过读写分离实现并发:
写操作:仅更新可见性位图(内存操作),无锁 读操作:基于时间戳生成可见性位图,无锁 冲突处理:通过时间戳排序,写入时保证原子性
Q2: 宕机后如何恢复?
A: 三层恢复机制:
检查点(Checkpoint) :定期将可见性状态写入磁盘 写缓冲持久化:MemTable 数据先刷到对象存储 Geddon WAL:Storage GC 的恢复日志
Q3: 如何与现有查询引擎集成?
A: 通过 RetinaServer 的 gRPC 接口:
serviceRetinaWorkerService{
rpc AddVisibility(AddVisibilityRequest) returns (AddVisibilityResponse);
rpc QueryVisibility(QueryVisibilityRequest) returns (QueryVisibilityResponse);
rpc GetWriteBuffer(GetWriteBufferRequest) returns (GetWriteBufferResponse);
rpc UpdateRecord(UpdateRecordRequest) returns (UpdateRecordResponse);
}
Trino 等引擎通过 Pixels Connector 调用这些接口,实现:
写入时更新可见性 读取时过滤已删除行
Q4: 为什么选择 C++ 实现可见性层?
A: 性能考量:
内存效率:C++ 位图操作更高效 低延迟:避免 JNI 调用开销(批量操作) SIMD 优化:C++ 可利用 SIMD 加速位图运算
八、总结
8.1 核心价值
Pixels-Retina 的核心价值可以用一句话概括:
"在数据湖上实现了数据库级别的实时变更能力,
同时保留了数据湖的大规模分析查询能力"
8.2 技术亮点
轻量级 MVCC:无需完整事务机制,通过删除链实现高效并发控制 VFoR (Vectorized-Filter-on-Read) :读取时向量化过滤,性能优异 无 MoR Merge:写入即可见,消除读放大问题 C++ 高性能层:可见性位图操作使用原生代码,延迟极低
8.3 适用边界
适合:
需要 CDC 实时同步的场景 既有高频写入又有大规模查询的场景 对数据新鲜度有严格要求(<1秒)
不适合:
强事务要求(ACID)的场景 需要复杂查询(DQL+DML 混合)的场景 超大规模单表(>10TB)的纯分析场景
参考
Pixels 论文 (ICDE'25, SIGMOD'23, ICDE'22)
DeepWiki Architecture
https://github.com/pixelsdb/pixels