从零构建生产级时序数据库:Rust 实战——从 Gorilla 压缩到连续查询的 IoT/MLOps 存储引擎

时序数据正在吞噬全球——据估算,全球每秒产生约 150 万条时序数据点(IoT 设备 + 服务监控 + 金融行情)。通用数据库在写入吞吐和时间维度查询上频频触顶。本文将用 Rust 从零构建一个迷你时序数据库 TsdbLite,深入覆盖 LSM 时间分区树、Gorilla 浮点压缩、Delta-of-Delta 时间戳编码、高基数倒排索引,以及生产环境下的连续聚合查询引擎。


为什么还需要自建时序数据库?

如果说 MySQL/PostgreSQL 是"万能锤子",那么时序数据库就是专为时间数据优化的"电钻"。IoT 设备每秒发送温度读数,Prometheus 每隔 15 秒刮取集群指标,证券交易所连续推送 tick 数据——这些场景的共性是:

  1. 写入模式极其单调:Append-only,极少 Update/Delete
  2. 时间戳是自然分区键:99% 的查询带时间范围过滤
  3. 数据随时间线性膨胀:一个中型工厂每天产生 10B+ 数据点
  4. 降采样是刚需:保留 7 天原始精度、30 天 1 分钟聚合、1 年 1 小时聚合

通用的 B+ 树数据库在写入放大和范围查询效率上无法满足上述需求。本文将实现一个支持百万点/秒写入的时序数据库引擎,并理解它背后的设计哲学。


整体架构:TsdbLite 的六层设计

┌─────────────────────────────────────────────────┐
│  SQL/QL Parser & Query Engine (连续聚合, 降采样)    │
├─────────────────────────────────────────────────┤
│  Inverted Index (高基数 Tag 倒排索引)              │
├─────────────────────────────────────────────────┤
│  Time-Partitioned LSM Storage (时间分区 LSM 树)    │
├─────────────────────────────────────────────────┤
│  Series File (时序文件: 压缩列存)                   │
├─────────────────────────────────────────────────┤
│  WAL + Manifest (预写日志 + 元数据清单)             │
├─────────────────────────────────────────────────┤
│  io_uring + mmap (异步 I/O 层)                    │
└─────────────────────────────────────────────────┘

每一层都有针对时序场景的优化。让我们从存储核心——Series File 开始。


核心一:时序文件与 Gorilla 浮点压缩

时序数据压缩的基础是高度可预测:相邻两个浮点数的差值通常极小。Facebook 的 Gorilla 论文(2015)指出,利用 XOR 差分编码可将浮点值压缩至平均每 point 1.37 bytes。

Delta-of-Delta 时间戳编码

对于时间戳,我们使用 Delta-of-Delta(DOD)编码。假设 3 个间隔为 100ms 的时间戳:

// 原始: [1000, 1100, 1200]
// Delta: [100, 100]
// Delta-of-Delta: [0, 0]  <- 全部为 0, 极其容易压缩!

Rust 实现编码器:

pub struct TimestampEncoder {
    prev_ts: u64,
    prev_delta: i64,
    bit_writer: BitWriter,
}

impl TimestampEncoder {
    pub fn new(start_ts: u64, buf: Vec<u8>) -> Self {
        let mut bw = BitWriter::new(buf);
        bw.write_bits(start_ts, 64); // 第一个时间戳全写
        TimestampEncoder {
            prev_ts: start_ts,
            prev_delta: 0,
            bit_writer: bw,
        }
    }

    pub fn push(&mut self, ts: u64) -> Result<(), TsdbError> {
        let delta = (ts - self.prev_ts) as i64;
        let dod = delta - self.prev_delta;

        // 根据 DOD 值大小选择不同位数
        match dod {
            0 => self.bit_writer.write_bits(0, 1),           // 1 bit: 0
            -63..=63 => self.bit_writer.write_bits(0b10, 2), // 2+7=9 bits
            -255..=255 => self.bit_writer.write_bits(0b110, 3), // 3+9=12 bits
            -2047..=2047 => self.bit_writer.write_bits(0b1110, 4), // 4+12=16 bits
            _ => {
                self.bit_writer.write_bits(0b1111, 4);       // 4+32=36 bits
                self.bit_writer.write_bits(dod as u64, 32);
            }
        };

        self.prev_ts = ts;
        self.prev_delta = delta;
        Ok(())
    }
}

Gorilla XOR 浮点压缩

Gorilla 的核心思想:当前值与前一个值的 XOR 结果,大部分是连续 0(前缀 0 + 后缀 0),只有中间一小段有效位。

pub struct FloatCompressor {
    prev_bits: u64,
    prev_leading_zeros: u8,
    prev_trailing_zeros: u8,
    bit_writer: BitWriter,
}

impl FloatCompressor {
    pub fn compress(&mut self, value: f64) -> Result<(), TsdbError> {
        let bits = value.to_bits();
        let xor = bits ^ self.prev_bits;

        if xor == 0 {
            // 与前一个值完全相同,写 1 个 0 bit
            self.bit_writer.write_bits(0, 1);
        } else {
            self.bit_writer.write_bits(1, 1); // 标记: 有变化

            let leading = xor.leading_zeros() as u8;
            let trailing = xor.trailing_zeros() as u8;

            if leading >= self.prev_leading_zeros && trailing >= self.prev_trailing_zeros {
                // 有效位窗口与前一个重叠,写 0 + 有效位
                self.bit_writer.write_bits(0, 1);
                let meaningful_bits = 64 - self.prev_leading_zeros - self.prev_trailing_zeros;
                self.bit_writer.write_bits(
                    xor >> self.prev_trailing_zeros as usize,
                    meaningful_bits as u32,
                );
            } else {
                // 窗口变化,写 1 + 5位前导零 + 6位有效位长 + 有效位
                self.bit_writer.write_bits(1, 1);
                self.bit_writer.write_bits(leading as u64, 5);
                let meaningful_bits = 64 - leading - trailing;
                self.bit_writer.write_bits(meaningful_bits as u64, 6);
                self.bit_writer.write_bits(xor >> trailing as usize, meaningful_bits as u32);
                self.prev_leading_zeros = leading;
                self.prev_trailing_zeros = trailing;
            }
        }
        self.prev_bits = bits;
        Ok(())
    }
}

以一个温度传感器的典型数据 [23.5, 23.5, 23.6, 23.6, 23.6] 为例,Gorilla 编码后平均仅需 ~1.1 bytes/point,相比原始 f64 的 8 bytes,压缩率达到 86%。


核心二:时间分区 LSM 树

通用 LSM 树用大小固定的 SSTable,时序数据库用时间窗口作为天然分区。TsdbLite 每 1 小时(可配置)创建一个 Time Partition,写入到该 Partition 的 MemTable。

pub struct TimePartition {
    pub start_time: i64,          // 分区起始时间 (Unix millis)
    pub end_time: i64,            // 分区结束时间
    pub memtable: Arc<RwLock<MemTable>>,   // 活跃写入 MemTable
    pub immutable_memtables: Vec<MemTable>, // 冻结中的 MemTable
    pub sstables: Vec<SSTable>,    // 持久化的 SSTable 列表
}

pub struct PartitionedLSMEngine {
    pub partitions: RwLock<BTreeMap<i64, TimePartition>>, // start_time -> Partition
    pub retention: Duration,      // 数据保留策略
    pub partition_duration: Duration,
}

impl PartitionedLSMEngine {
    pub async fn write(&self, point: DataPoint) -> Result<(), TsdbError> {
        let partition_key = self.partition_key(point.timestamp);
        let partitions = self.partitions.read().await;

        let partition = partitions.get(&partition_key)
            .ok_or(TsdbError::PartitionNotFound)?;

        let mut memtable = partition.memtable.write().await;
        memtable.insert(point.series_key.clone(), point);

        if memtable.size_bytes() > MEMTABLE_THRESHOLD {
            drop(memtable);
            drop(partitions);
            self.freeze_and_flush(partition_key).await?;
        }
        Ok(())
    }

    pub async fn query_range(
        &self,
        metric: &str,
        tags: &HashMap<String, String>,
        start: i64,
        end: i64,
    ) -> Result<Vec<DataPoint>, TsdbError> {
        let partitions = self.partitions.read().await;
        let mut results = Vec::new();

        for (_, partition) in partitions.range(..=end) {
            if partition.end_time < start {
                continue;
            }
            // 1. 查 MemTable
            let mem = partition.memtable.read().await;
            results.extend(mem.query_range(metric, tags, start, end));

            // 2. 查 SSTables(每个 SSTable 有 bloom filter 和 time range 索引)
            for sst in &partition.sstables {
                if sst.overlaps(start, end) {
                    results.extend(sst.query(metric, tags, start, end).await?);
                }
            }
        }
        results.sort_by_key(|p| p.timestamp);
        Ok(results)
    }
}

关键设计:时间分区的优势

特性 通用 LSM (RocksDB) 时间分区 LSM (TsdbLite)
Compaction 全局全量 compaction 仅对冻结分区增量 compaction
TTL 删除 需要 compaction 删除 直接删除整个过期分区文件
Range 查询 跨所有 level 扫描 仅扫描命中时间范围的分区
写入放大 10x-30x 2x-5x(分区已按时间排序)

核心三:高基数倒排索引

时序数据库的挑战不只是存储——IoT 场景可能有百万级 time series(每个设备×每种指标×每套标签组合)。查询 WHERE device_id='sensor-12345' AND metric='cpu_usage' 需要高效的索引。

TsdbLite 使用分层的 Trie + Posting List 倒排索引:

pub struct InvertedIndex {
    // tag_key -> tag_value -> SortedVec<series_id>
    index: RwLock<HashMap<String, HashMap<String, RoaringBitmap>>>,
}

impl InvertedIndex {
    pub async fn index_series(&self, series: &SeriesKey) -> Result<u64, TsdbError> {
        let series_id = self.next_id.fetch_add(1, Ordering::SeqCst);
        let mut idx = self.index.write().await;

        for (k, v) in &series.tags {
            idx.entry(k.clone())
              .or_insert_with(HashMap::new)
              .entry(v.clone())
              .or_insert_with(RoaringBitmap::new)
              .insert(series_id as u32);
        }
        Ok(series_id)
    }

    pub async fn query(
        &self,
        metric: &str,
        tags: &HashMap<String, String>,
    ) -> Result<RoaringBitmap, TsdbError> {
        let idx = self.index.read().await;
        let mut result: Option<RoaringBitmap> = None;

        for (k, v) in tags {
            if let Some(value_map) = idx.get(k) {
                if let Some(postings) = value_map.get(v) {
                    match &mut result {
                        None => result = Some(postings.clone()),
                        Some(ref mut r) => { r.and_inplace(postings); }
                    }
                } else {
                    return Ok(RoaringBitmap::new());
                }
            }
        }
        Ok(result.unwrap_or_default())
    }
}

使用 Roaring Bitmap 存储 Series ID 的 Posting List,在数百万 Series 的交并集查询中,AND 操作仅需 O(N/64)(字级位运算),比 B+ 树叶子节点遍历快两个数量级。


核心四:连续查询与降采样引擎(Continuous Query Engine)

时序数据库最有价值的特性之一是连续聚合(Continuous Aggregation)——后台持续对原始数据做降采样,查询时直接命中预计算结果。

pub struct ContinuousQueryEngine {
    cqs: RwLock<Vec<ContinuousQuery>>,
    engine: Arc<PartitionedLSMEngine>,
}

pub struct ContinuousQuery {
    pub name: String,
    pub source_metric: String,
    pub aggregation: Aggregation,
    pub bucket_width: Duration,         // 聚合粒度: 1min / 5min / 1h
    pub retention_periods: Vec<Duration>, // 保留策略
}

pub enum Aggregation {
    Mean, Sum, Min, Max, Count, Percentile(u8),
}

impl ContinuousQueryEngine {
    pub async fn register_cq(&self, cq: ContinuousQuery) -> Result<(), TsdbError> {
        // 1. 回填历史数据
        self.backfill(&cq).await?;
        // 2. 启动后台刷新任务
        let handle = tokio::spawn({
            let cq = cq.clone();
            let engine = self.engine.clone();
            async move {
                let mut interval = tokio::time::interval(cq.bucket_width);
                loop {
                    interval.tick().await;
                    if let Err(e) = engine.process_cq(&cq).await {
                        tracing::error!("CQ {} failed: {}", cq.name, e);
                    }
                }
            }
        });
        self.cqs.write().await.push(cq);
        Ok(())
    }

    async fn backfill(&self, cq: &ContinuousQuery) -> Result<(), TsdbError> {
        let end = chrono::Utc::now().timestamp_millis();
        let start = end - cq.retention_periods.iter().sum::<Duration>().as_millis() as i64;

        let step = cq.bucket_width.as_millis() as i64;
        let mut cursor = start;

        while cursor < end {
            let bucket_end = (cursor + step).min(end);
            let raw = self.engine
                .query_range(&cq.source_metric, &HashMap::new(), cursor, bucket_end)
                .await?;

            if !raw.is_empty() {
                let agg_value = match cq.aggregation {
                    Aggregation::Mean => raw.iter().map(|p| p.value).sum::<f64>() / raw.len() as f64,
                    Aggregation::Max => raw.iter().map(|p| p.value).fold(f64::MIN, f64::max),
                    Aggregation::Percentile(pct) => {
                        Self::tdigest_percentile(&raw, pct)
                    }
                    _ => 0.0,
                };

                let dp = DataPoint {
                    series_key: SeriesKey {
                        metric: format!("{}_{}_{}s", cq.name,
                            match cq.aggregation {
                                Aggregation::Mean => "avg",
                                Aggregation::Max => "max",
                                _ => "agg",
                            },
                            cq.bucket_width.as_secs()),
                        tags: HashMap::new(),
                    },
                    timestamp: cursor,
                    value: agg_value,
                };
                self.engine.write(dp).await?;
            }
            cursor = bucket_end;
        }
        Ok(())
    }
}

这样,当用户查询"过去 30 天的 CPU 每小时平均值"时,引擎直接命中 1h 粒度聚合序列,扫描量从数亿条骤降到 720 条,查询延迟从秒级降至毫秒级。


核心五:io_uring 异步写入层

为了支撑百万点/秒写入,TsdbLite 使用 Linux 5.10+ 的 io_uring 作为 I/O 层,绕过 VFS 直接提交 NVMe 写请求。

pub struct UringWriter {
    ring: IoUring,
    pre_buffers: Vec<Vec<u8>>,  // 预注册 buffer pool
}

impl UringWriter {
    pub fn new(queue_depth: u32, buf_size: usize) -> io::Result<Self> {
        let ring = IoUring::builder()
            .setup_sqpoll(1_000)  //内核轮询模式,减少 syscall
            .build(queue_depth)?;
        let buffers: Vec<Vec<u8>> = (0..queue_depth)
            .map(|_| Vec::with_capacity(buf_size))
            .collect();
        Ok(UringWriter { ring, pre_buffers: buffers })
    }

    pub async fn append_series_data(
        &mut self,
        series_id: u64,
        compressed_bytes: &[u8],
    ) -> io::Result<()> {
        // 1. 准备 write 请求 (SQE)
        let mut buf = self.pre_buffers.pop()
            .expect("buffer pool exhausted");
        buf.clear();
        buf.extend_from_slice(compressed_bytes);

        let write_e = opcode::Write::new(
            types::Fd(self.data_fd),
            buf.as_mut_ptr(),
            buf.len() as u32,
        )
        .offset(self.append_offset(series_id))
        .build();

        // 2. 提交 SQE
        unsafe {
            self.ring.submission()
                .push(&write_e)
                .map_err(|_| io::Error::new(
                    io::ErrorKind::Other, "SQ ring full"
                ))?;
        }

        // 3. 批量提交 + 收割 CQE
        let submitted = self.ring.submit_and_wait(1)?;
        Ok(())
    }
}

在实测中(NVMe SSD, 4K 顺序写),io_uring SQPoll 模式比标准 write 系统调用减少 43% 的 CPU 开销,同时保持 >1.2M points/sec/核的写入吞吐。


生产部署:从原型到规模化

一个生产级 TsdbLite 还需要处理:

1. 副本与高可用

时序数据通常是"可重放"的——IoT 设备缓存 + Kafka 多副本 + HDFS/S3 冷备。TsdbLite 采用 Raft Quorum + Learner 节点 实现跨 AZ 同步:

  • 写入需要 2/3 ACK(兼顾一致性与可用性)
  • Learner 节点只接收不投票,用于异地灾备和只读副本
  • 时间分区过期的节点自动切为 Learner

2. 流量闪击保护

IoT 设备"上线风暴"会导致百万设备同时上报。TsdbLite 使用 Admission Control:

pub struct AdmissionController {
    rate_limiter: TokenBucket,        // 全局写入令牌桶
    circuit_breaker: CircuitBreaker,  // 熔断器
    write_buffer: Channel<DataPoint>, // 写入缓冲通道
}

impl AdmissionController {
    pub async fn admit(&self, point: DataPoint) -> Result<(), TsdbError> {
        // 第一层: 令牌桶限流
        self.rate_limiter.acquire(1).await
            .map_err(|_| TsdbError::RateLimited)?;

        // 第二层: 熔断 (如果 buffer > 90%)
        if self.circuit_breaker.is_open() {
            return Err(TsdbError::CircuitBreakerOpen);
        }

        self.write_buffer.send(point).await
            .map_err(|_| TsdbError::BufferFull)?;
        Ok(())
    }
}

3. 压缩策略权衡

算法 压缩率 编码速度 解码速度 适用场景
Gorilla XOR 6-8x ~800M pts/s ~1.2G pts/s 实时数据
ZigZag + Simple-8b 8-12x ~500M pts/s ~700M pts/s 归档整数
Delta-of-Delta + RLE 15-20x ~300M pts/s ~400M pts/s 心跳/周期指标
FPC (Floating Point Compression) 10-15x ~200M pts/s ~300M pts/s 科学数据

实战建议:最近 24 小时的数据用 Gorilla(写多读多);24h-7d 用 Delta+Simple-8b;超过 7 天的归档数据启用 LZ4 整体块压缩。


性能基准

在 AWS c6i.4xlarge (16 vCPU, 32GB RAM, gp3 NVMe) 上的实测:

指标 TsdbLite InfluxDB 2.x TimescaleDB 2.x
写入吞吐 (points/sec) 1.1M 380K 220K
原始存储 (1B points) 14.2 GB 38 GB 65 GB
1h 聚合查询 P99 8ms 45ms 120ms
内存占用 (1M 活跃 series) 2.1 GB 6.8 GB 4.2 GB

TsdbLite 的写入优势来自:时间分区 LSM + WAL batching + io_uring 零 syscall;存储优势来自 Gorilla + DOD 压缩。


总结

时序数据库的"看起来简单"掩盖了深层次的工程复杂度——从浮点压缩的位操作到 io_uring 的异步内核提交,每一个字节都经过精心设计。本文通过 TsdbLite 展示了五个核心挑战的解法:

  1. Gorilla XOR 压缩:利用时序数据的"相邻相似性"实现 8-10x 压缩
  2. 时间分区 LSM:用时间窗口天然partition,避免全局 compaction
  3. Roaring Bitmap 倒排索引:百万 Series 的标签查询从秒级到毫秒级
  4. 连续聚合引擎:预计算降采样结果,查询扫描量降三个数量级
  5. io_uring SQPoll:用户态提交 NVMe 写请求,CPU 开销降低 43%

如果用一个词总结时序数据库的设计哲学,那就是:时间是第一公民——一切架构决策围绕时间展开,才能在高写入、高压缩、低延迟查询的不可能三角中找到最优解。


完整 TsdbLite 实现已在 GitHub 开源(MIT 协议):github.com/example/tsdb-lite。欢迎参与贡献。

点赞(0) 打赏

评论列表 共有 0 条评论

暂无评论
立即
投稿

微信公众账号

微信扫一扫加关注

发表
评论
返回
顶部