物流轨迹的时序存储:亿级包裹的实时位置查询与历史轨迹回放
物流轨迹的时序存储亿级包裹的实时位置查询与历史轨迹回放一、快递到哪了背后的查询洪水双11凌晨某快递公司的物流查询接口QPS从日常的5万飙升至300万。用户反复刷新查询同一个包裹的物流轨迹——一个典型的读多写少场景。但问题在于每个查询都需要按时间顺序展示该包裹的20-30条轨迹节点揽件→中转→派送→签收这意味着每次查询都是一次ORDER BY scan_time DESC LIMIT 30。当30亿包裹每个包裹30条轨迹的查询同时涌来时传统的MySQL在(waybill_no, scan_time)联合索引下需要先定位到waybill_no的数据页再在页内按时间排序。300万QPS下Buffer Pool被冲刷得一塌糊涂。二、轨迹数据的时序存储方案三、ClickHouse轨迹表与查询实现CREATE TABLE waybill_trajectories ON CLUSTER logistics ( waybill_no String, scan_time DateTime64(3), scan_type LowCardinality(String), -- COLLECT/TRANSIT/DELIVERY/SIGN location_code String, -- 网点/中转场编码 location_name String, city LowCardinality(String), province LowCardinality(String), courier_id String, latitude Float64, longitude Float64, signer_name String, exception_code String, -- 异常类型延误/破损/退回 create_time DateTime DEFAULT now() ) ENGINE ReplicatedReplacingMergeTree( /clickhouse/tables/{shard}/waybill_trajectories, {replica}, create_time ) PARTITION BY toYYYYMM(scan_time) ORDER BY (waybill_no, scan_time) TTL scan_time INTERVAL 30 DAY TO VOLUME hot_ssd, scan_time INTERVAL 90 DAY TO VOLUME cold_s3 SETTINGS index_granularity 8192;查询与实时追踪实现class WaybillTracker: def __init__(self, clickhouse_client, mysql_pool, redis_client): self.ch clickhouse_client self.mysql mysql_pool self.redis redis_client def get_trajectory(self, waybill_no: str, days: int 30) - dict: 查询包裹轨迹 cache_key ftrajectory:{waybill_no} # L1: Redis缓存高频查询包裹 try: cached self.redis.get(cache_key) if cached: trajectory json.loads(cached) # 检查是否有新的扫描记录 last_scan trajectory[-1][scan_time] if trajectory else latest self._get_latest_scan(waybill_no, last_scan) if latest: trajectory.extend(latest) self.redis.setex(cache_key, 300, json.dumps(trajectory)) return {waybill_no: waybill_no, traces: trajectory} except RedisError: pass # L2: ClickHouse最近30天 try: trajectory self._query_clickhouse(waybill_no, days) if trajectory: self.redis.setex(cache_key, 300, json.dumps(trajectory)) return {waybill_no: waybill_no, traces: trajectory} except Exception: pass # L3: MySQL历史数据 return self._query_mysql(waybill_no) def _query_clickhouse(self, waybill_no: str, days: int) - list: ClickHouse查询轨迹 query SELECT scan_time, scan_type, location_name, city, province, courier_id, latitude, longitude, exception_code FROM waybill_trajectories WHERE waybill_no %(wn)s AND scan_time now() - INTERVAL %(days)s DAY ORDER BY scan_time ASC try: result self.ch.execute(query, { wn: waybill_no, days: days }) return [ { scan_time: str(row[0]), scan_type: row[1], location: row[2], city: row[3], province: row[4], courier: row[5], lat: row[6], lng: row[7], exception: row[8] } for row in result ] except Exception as e: raise TrajectoryQueryException(f轨迹查询失败: {waybill_no}, e) def scan_event(self, waybill_no: str, event: dict): 处理实时扫描事件 # 写入ClickHouse try: self.ch.execute( INSERT INTO waybill_trajectories VALUES, [(waybill_no, datetime.now(), event[type], event[location_code], event[location_name], event[city], event[province], event[courier_id], event[lat], event[lng], , event.get(exception, ), datetime.now())] ) except Exception as e: raise ScanEventException(f扫描事件写入失败: {waybill_no}, e) # 清除Redis缓存强制下次查询走ClickHouse try: self.redis.delete(ftrajectory:{waybill_no}) except RedisError: pass # 检测异常事件 if event.get(exception): self._handle_exception(waybill_no, event) def get_active_waybills_by_city(self, city: str) - dict: 实时查询某城市的活跃包裹数 query SELECT count(DISTINCT waybill_no) AS active_count, countIf(exception_code ! ) AS exception_count, -- 按扫描类型分布 countIf(scan_type DELIVERY) AS delivering, countIf(scan_type SIGN) AS signed_today FROM waybill_trajectories WHERE city %(city)s AND scan_time today() GROUP BY city result self.ch.execute(query, {city: city}) if result: row result[0] return { city: city, active_waybills: row[0], exceptions: row[1], delivering: row[2], signed_today: row[3] } return {city: city, active_waybills: 0}四、轨迹时序存储的四个关键设计关键一查询模式决定排序键。99%的查询是WHERE waybill_no ? ORDER BY scan_time所以排序键必须是(waybill_no, scan_time)。ClickHouse按排序键物理排序存储这个查询的扫描效率最高。关键二异常轨迹的实时告警。包裹在某中转场停留超过24小时→可能丢件。需要有旁路的Flink作业持续监控最近一条扫描记录的scan_type ! SIGN AND scan_time now() - 24h触发告警。关键三数据归档的查询可用性。30天以上的轨迹迁移到S3后用户查询历史包裹时延迟从100ms升到2秒。需要在UI上明确标注历史轨迹查询可能较慢并设置30秒超时。关键四隐私数据的生命周期。签收人姓名、电话号码在物流轨迹中属于个人信息。签收后7天signer_name字段应自动脱敏仅保留姓氏30天后从热数据中移除。五、总结物流轨迹是时序数据的典型场景写入是追加式的高吞吐每秒万条扫描事件查询是点查式的高并发单个包裹的轨迹存储是滚动式的冷热分离30天热、90天温、之后归档。ClickHouse的MergeTree引擎在ORDER BY TTL的组合下完美适配这个场景。一句大白话总结轨迹数据就像体温计记录的是某个时刻发生了什么而不是当前状态是什么。时序数据库天然适合这种数据形态。本文属于「行业场景与项目复盘」系列深入分析物流轨迹的时序数据存储与实时查询方案。