PostgreSQL码农集散地

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) 机制,存在以下问题:

问题
描述
数据新鲜度差
分钟级延迟,因为依赖批量快照合并
变更处理吞吐低
数据变更需要重新写入整个文件
查询性能不稳定
Merge 操作消耗大量计算资源
资源成本高
需要频繁的 compaction 和版本合并

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 版本数据                                        │  
│  • 基于时间戳的可见性判断                                         │  
│  • 无锁读,仅在写入时轻量同步                                      │  
└─────────────────────────────────────────────────────────────────┘  

核心数据结构:

  1. VersionedData:存储基础位图和基准时间戳

    structVersionedData {
    long* baseBitmap;      // 基础可见性位图  
    long baseTimestamp;    // 基准时间戳  
    };  
  2. DeleteIndexBlock:删除链节点

    structDeleteIndexBlock {
    uint64_t* entries;    // 打包的 (rowId, timestamp) 对  
        DeleteIndexBlock* next;  
    };  
  3. 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-RetinaApache IcebergApache Paimon
数据新鲜度
10ms 级
分钟级
分钟级
变更粒度
行级
文件级
行级
MVCC 机制
轻量级删除链
快照隔离
快照隔离
合并策略
无需 MoR
Merge-On-Read
Merge-On-Read
索引支持
多版本主索引
无内置
主键索引
存储格式
自研 Pixels 列存
Parquet/ORC
Parquet/ORC
CDC 支持
原生 CDC 重放
需额外组件
原生支持
查询引擎
Trino/DuckDB 等
通用
通用

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 性能数据对比

根据项目文档和论文:

指标
Pixels-Retina
Iceberg
Paimon
数据新鲜度
10ms
分钟级
分钟级
变更重放吞吐
3.2M rows/s
~100K rows/s
~200K rows/s
查询延迟
相当
相当
相当
存储成本
较低(无版本堆积)
中等
中等

五、使用前后效果变化

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  

部署说明

  1. pixels-sink 是一个独立的 Java 进程(不是 Flink Job)
  2. 它订阅 Kafka 中 Debezium 的 topic
  3. 消费 CDC 消息,重建为镜像事务
  4. 通过 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 常见问题处理

问题
原因
解决方案
数据延迟增加
GC 积压
减小 GC 间隔,增加 GC 线程
内存持续增长
删除链未压缩
检查 safeGcTimestamp 是否正常推进
写入 QPS 下降
MemTable 竞争
增加 virtual.node.num
查询偶尔超时
版本切换
检查 SuperVersion 锁竞争

七、关键技术问题解答

基于 DeepWiki 架构分析,关于 Retina 的关键问题:

Q1: Retina 如何保证读写并发?

A: Retina 通过读写分离实现并发:

  • 写操作:仅更新可见性位图(内存操作),无锁
  • 读操作:基于时间戳生成可见性位图,无锁
  • 冲突处理:通过时间戳排序,写入时保证原子性

Q2: 宕机后如何恢复?

A: 三层恢复机制:

  1. 检查点(Checkpoint) :定期将可见性状态写入磁盘
  2. 写缓冲持久化:MemTable 数据先刷到对象存储
  3. 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 技术亮点

  1. 轻量级 MVCC:无需完整事务机制,通过删除链实现高效并发控制
  2. VFoR (Vectorized-Filter-on-Read) :读取时向量化过滤,性能优异
  3. 无 MoR Merge:写入即可见,消除读放大问题
  4. C++ 高性能层:可见性位图操作使用原生代码,延迟极低

8.3 适用边界

适合:

  • 需要 CDC 实时同步的场景
  • 既有高频写入又有大规模查询的场景
  • 对数据新鲜度有严格要求(<1秒)

不适合:

  • 强事务要求(ACID)的场景
  • 需要复杂查询(DQL+DML 混合)的场景
  • 超大规模单表(>10TB)的纯分析场景

参考

Pixels 论文 (ICDE'25, SIGMOD'23, ICDE'22)

DeepWiki Architecture

https://github.com/pixelsdb/pixels