Milvus 混合时间戳(Hybrid TSO)全解析:位布局、分配机制与 UTC 时间还原
【免费下载链接】milvusMilvus is a high-performance, cloud-native vector database built for scalable vector ANN search项目地址: https://gitcode.com/GitHub_Trending/mi/milvus
Milvus 是一套云原生分布式向量数据库,其所有写入与删除等 DML 事件都依赖统一的全局时间戳(TSO)来保证跨节点的一致性可见性。本文以设计文档 20211214-milvus_hybrid_ts.md 为核心,结合仓库中internal/tso、pkg/util/tsoutil的真实实现,系统讲解 Hybrid TSO 的二进制位布局、RootCoord 侧的分配流水线,以及如何把uint64时间戳还原成可读的 UTC 时间。读完本文,你将掌握 Milvus 时间戳的编码格式、官方提供的解析工具函数,并能独立写出 Go / Python 环境下的解析代码。
背景:为什么 Milvus 需要 TSO
Milvus 是分布式系统,客户端可能连接任意一个 Proxy 节点。如果每个节点都使用本地时钟为事件打时间戳,就会面临两类典型问题(更完整的论述见配套设计文档 20211215-milvus_timesync.md):
- 时钟不同步:不同物理机的本地时钟存在偏差,极端情况下某一节点的事件会被另一节点"滞后一天"看到;
- 网络延迟:某个操作在时间
t17发起,但早于t17的事件(如 Delete)仍滞留在消息队列中,会破坏快照读的一致性。
解决方案是所有事件的时间戳不再取自本地时钟,而是统一向Timestamp Oracle(TSO 服务)申请。该服务由 RootCoord 承担,对外暴露AllocTimestamp/AllocTimestampResponse等 RPC,而 TSO 本身采用的是与 TiKV 一脉相承的实现思路。此外,Milvus 与 TSO 配套还有一套基于 TimeTick 消息的"时间同步系统"(各 Proxy 周期上报各 channel 的最新时间戳,RootCoord 取最小值回插到消息流),用于解决"读到某个消息时如何确认更小时间戳的消息都已消费完毕"这一一致性问题。
用户(如 DBA)有时希望把分布式系统里分散在各节点的事件按真实发生的 UTC 时间排序。由于事件的时间戳都来自同一个 TSO,事件之间的先后关系可以直接用 TSO 大小比较;因此剩下要回答的问题就是——给定一个uint64的 TSO,如何还原出它的 UTC 物理时间?
Hybrid TSO 的位布局
Milvus 中 TSO 的底层类型是uint64,由两部分拼装而成:
| 位域 | 长度 | 含义 |
|---|---|---|
| physical(高位) | 前 46 bits | 物理部分,即 UTC 时间,单位为毫秒 |
| logical(低位) | 后 18 bits | 逻辑部分,用于同一毫秒内区分多次分配 |
这一 46/18 划分与仓库源码保持一致:pkg/util/tsoutil/tso.go中定义了logicalBits = 18与logicalBitsMask = (1 << logicalBits) - 1,并据此提供合成与解析工具函数:
ComposeTS(physical, logical) uint64:把物理毫秒值左移 18 位再叠加逻辑值,完成编码;ParseTS(ts) (time.Time, uint64):还原出time.Time与逻辑值;ParseHybridTs(ts) (int64, int64):还原出(物理毫秒, 逻辑值)两个整数。
从位数上看,物理部分可表达的上限约为2^46毫秒,足以覆盖当前 epoch 毫秒量级的全部应用场景;而低 18 位意味着在同一毫秒内最多可区分2^18 ≈ 262144个递增时间戳。当逻辑部分在单个毫秒内被"用尽"时,分配器会强制将物理时间向前推进一个毫秒,保证时间戳严格单调递增。
分配侧的位数约束
逻辑位宽的上限约束在分配器一侧同样得到呼应。internal/tso/tso.go定义了:
const ( // UpdateTimestampStep is used to update timestamp. UpdateTimestampStep = 50 * time.Millisecond // updateTimestampGuard is the min timestamp interval. updateTimestampGuard = time.Millisecond // maxLogical is the max upper limit for logical time. // When a TSO's logical time reaches this limit, // the physical time will be forced to increase. maxLogical = int64(1 << 18) )其中maxLogical = 1 << 18正是低 18 位的容量上限。分配器在GenerateTSO中通过atomic.AddInt64自增逻辑部分;一旦logical >= maxLogical且启用了LimitMaxLogic,就会先休眠UpdateTimestampStep(50ms)等待物理时间窗口前移后重试,从而规避同一毫秒内逻辑位溢出的风险。内存中的当前时间窗口保存在atomicObject{physical time.Time; logical int64}结构里,并通过unsafe.Pointer做无锁读写。
TSO 在 RootCoord 侧如何被分配
TSO 的分配集中在 RootCoord。internal/tso/global_allocator.go中的GlobalTSOAllocator实现了Allocator接口,其内部持有一个timestampOracle:
// GenerateTSO is used to generate a given number of TSOs. func (gta *GlobalTSOAllocator) GenerateTSO(count uint32) (uint64, error) { ... for i := 0; i < maxRetryCount; i++ { current := (*atomicObject)(atomic.LoadPointer(>a.tso.TSO)) ... physical = current.physical.UnixMilli() logical = atomic.AddInt64(¤t.logical, int64(count)) if logical >= maxLogical && gta.LimitMaxLogic { ... time.Sleep(UpdateTimestampStep) continue } return tsoutil.ComposeTS(physical, logical), nil } return 0, merr.WrapErrServiceInternalMsg("can not get timestamp") }整体分配流水线呈现为"定期推进物理时间窗口 + 内存内自增逻辑位"的组合:
- 初始化:
InitTimestamp会先从底层TxnKV(etcd/TiKV)读回上次持久化的时间窗口(以 Unix 纳秒的uint64大端字节存储,见saveTimestamp),避免重启后时间回退;若本机系统时间与已保存值之差小于updateTimestampGuard(1ms),则从已保存值上顺延 1ms 起步。 - 周期推进:
UpdateTimestamp默认每UpdateTimestampStep(50ms)被触发一次。若系统时间领先当前物理时间超过 1ms,则以系统时间为准;若prevLogical > maxLogical/2(逻辑部分消耗过半,说明单毫秒内分配过于密集),则把物理时间强制 +1ms 并清零逻辑位。持久化窗口(saveInterval = 3s)会被周期性地提前写入底层 KV,保证故障恢复后单调性不破。 - 即时分配:
GenerateTSO(count)读取当前物理毫秒,将logical原子自增count,再用tsoutil.ComposeTS拼出起始时间戳。Alloc(count)返回这批时间戳的起点,调用方按起点 + 偏移连续取值即可。
在实际调用链中,RootCoord 的各种 DDL / DML 元数据操作都依赖该分配器,例如 ddl_callbacks.go 中的ts, err := c.tsoAllocator.GenerateTSO(1),以及元数据表写入时批量申请时间戳的mt.tsoAllocator.GenerateTSO(2)等;对应的行为在 constraint_test.go 等测试中以 mock 断言的方式得到验证。
从 TSO 还原 UTC 时间(解析原理)
因为物理部分占据 TSO 的高位,且低位恰好是 18 bit 的逻辑值,所以还原 UTC 时间在原理上只需把 TSO 整体右移 18 位,拿到毫秒级 epoch 时间,再换算成人类可读的日期即可。设 TSO 类型为uint64:
physical(ms) = ts >> 18logical = ts & ((1 << 18) - 1)- 日期转换:
time.Unix(physical/1000, physical%1000 * 1e6)(秒 + 纳秒),或直接用 Python 的datetime.fromtimestamp(physical / 1000.0)。
ParseTS是仓库中唯一的官方 Go 实现,见 pkg/util/tsoutil/tso.go:
func ParseTS(ts uint64) (time.Time, uint64) { logical := ts & logicalBitsMask physical := ts >> logicalBits physicalTime := time.Unix(int64(physical/1000), int64(physical)%1000*time.Millisecond.Nanoseconds()) return physicalTime, logical }同文件中还提供若干便捷封装:
PhysicalTime(ts):只取物理time.Time;PhysicalTimeSeconds(ts):返回以秒计的浮点物理时间float64(ts>>logicalBits)/1000;ParseHybridTs(ts):返回(物理毫秒 int64, 逻辑值 int64)二元组,便于做时间差与日志输出;CalculateDuration(ts1, ts2):返回两时间戳物理部分相差的毫秒数;PhysicalTimeFormat(ts):直接格式化为"2006-01-02 15:04:05"字符串;Mod24H(ts)、AddPhysicalDurationOnTs(ts, duration)、SubByNow(ts):分别用于取"当日毫秒"、在物理部分上叠加时长、计算距今毫秒差等运维与监控场景;IsValidHybridTs/IsValidPhysicalTs:用于校验时间戳物理部分是否落在合法区间。
以上工具均配有单元测试 pkg/util/tsoutil/tso_test.go,例如Test_Tso验证了ComposeTSByTime→ParseHybridTs的往返一致性,TestCalculateDuration与TestAddPhysicalDurationOnTs验证了物理时间的差值计算与平移运算。
一个可验证的 Go 解析示例
把文档中的示例时间戳429164525386203142代入上述公式:
ts := uint64(429164525386203142) logical := ts & ((1 << 18) - 1) physical := ts >> 18 // 毫秒 fmt.Println(time.Unix(int64(physical/1000), int64(physical%1000)*int64(time.Millisecond)).UTC()) // 输出约为:2021-11-17 15:05:41 +0000 UTC也可以直接调用仓库现成的解析函数:
import "github.com/milvus-io/milvus/pkg/util/tsoutil" t, _ := tsoutil.ParseTS(429164525386203142) // t 为物理时间Python 环境下的时间戳解析
对不熟悉 Go 的开发者,设计文档给出了等价的 Python 实现(此处对注释做了补充说明):
>>> import datetime >>> LOGICAL_BITS = 18 >>> LOGICAL_BITS_MASK = (1 << LOGICAL_BITS) - 1 >>> def parse_ts(ts): ... logical = ts & LOGICAL_BITS_MASK ... physical = ts >> LOGICAL_BITS ... return physical, logical ... >>> ts = 429164525386203142 >>> utc_ts_in_milliseconds, _ = parse_ts(ts) >>> d = datetime.datetime.fromtimestamp(utc_ts_in_milliseconds / 1000.0) >>> d.strftime('%Y-%m-%d %H:%M:%S') '2021-11-17 15:05:41'需要注意:Python 的datetime.fromtimestamp默认返回本地时区时间而非严格意义上的 UTC。设计文档示例输出恰好与 UTC 对齐是因为运行环境处于 UTC 时区。若要得到确定性的 UTC 时间,建议显式指定时区:
>>> datetime.datetime.fromtimestamp(utc_ts_in_milliseconds / 1000.0, datetime.timezone.utc) datetime.datetime(2021, 11, 17, 15, 5, 41, tzinfo=datetime.timezone.utc)Hybrid TSO 的典型应用场景
- 事件排序与审计:同一集合上的 Insert / Delete 操作都携带单调递增的 TSO,直接按 TSO 数值排序即可得到与发生顺序一致的事件流;DBA 借助物理部分还原 UTC 时间后,可以按人类可读时间列出操作,便于审计与问题排查。
- 快照读 / 一致性时间推进:在 20211215-milvus_timesync.md 描述的机制中,各组件通过比较消息携带的 TSO 与已推进的 TimeTick 判断"是否可以安全消费",解析出物理时间也有助于监控系统时间水位与组件间延迟。
- 生命周期与过期判断:
SubByNow/CalculateDuration等工具函数常被用于计算数据距写入时刻的年龄、判断 TTL 与数据过期,是元数据与存储层的通用基础设施。
小结与延伸阅读
Milvus 的 Hybrid TSO 可以浓缩为三句话:
- 编码:
TSO = (物理毫秒 << 18) + 逻辑值,物理高 46 位、逻辑低 18 位,类型为uint64; - 分配:RootCoord 内的
GlobalTSOAllocator定期推进持久化时间窗口,在单毫秒内用原子自增填充逻辑位,任一环都保证单调递增; - 解析:
physical = ts >> 18即可还原 UTC 毫秒时间,官方 Go 封装位于tsoutil.ParseTS/ParseHybridTs。
如果想继续深入,推荐按以下顺序在仓库中阅读源码:
- 时间戳编解码与工具函数:pkg/util/tsoutil/tso.go 及其单元测试 pkg/util/tsoutil/tso_test.go;
- 全局分配器与
Allocator接口:internal/tso/global_allocator.go; - 时间窗口推进、持久化与复位逻辑:internal/tso/tso.go 及测试 internal/tso/global_allocator_test.go;
- 配套机制的整体设计:20211215-milvus_timesync.md。
结合本设计文档与上述实现代码,即可完整掌握 Milvus 分布式时间戳从分配、编码到解析的全链路。
【免费下载链接】milvusMilvus is a high-performance, cloud-native vector database built for scalable vector ANN search项目地址: https://gitcode.com/GitHub_Trending/mi/milvus
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考