ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

ClickHouse+S2索引实现万亿级轨迹数据秒级检索实战

ClickHouse+S2索引实现万亿级轨迹数据秒级检索实战 1. 项目概述当轨迹数据遇上“秒级”挑战在移动互联网和物联网时代轨迹数据正以前所未有的速度增长。从共享单车的骑行路线、物流车辆的配送轨迹到手机App的用户位置签到这些由时间、空间坐标构成的序列每天都在产生着海量记录。我们团队负责的“图灵平台”核心使命之一就是处理这类数据。当业务方提出“在数万亿条历史轨迹中快速找到经过某一片区域的所有车辆”这类需求时传统的数据库方案瞬间就“哑火”了。这不仅仅是数据量大的问题更是对查询响应时间的极限挑战——业务要求是“秒级”甚至“亚秒级”响应。这意味着从用户点击“查询”到看到结果整个链路的耗时必须控制在1秒以内背后涉及的是从存储引擎、索引设计到查询优化的全链路技术重构。今天我就结合我们平台从零到一构建万亿级轨迹数据秒级检索能力的实战经历拆解其中的核心思路、技术选型与那些“踩坑”得来的宝贵经验。2. 核心架构设计与技术选型背后的逻辑面对“万亿级”和“秒级”这两个关键词技术选型不能凭感觉必须基于数据特性和查询模式进行理性推演。2.1 为什么是ClickHouse在评估了HBase、Elasticsearch、Druid以及一些云原生数仓后我们最终将核心存储与计算引擎锚定在了ClickHouse上。这个决定基于几个关键考量首先查询模式匹配。我们的核心查询是典型的“点查”和“范围查”例如“查询某辆车在某个时间段内的轨迹”或“查询某个地理围栏内出现的所有车辆ID”。这类查询过滤条件强需要返回的数据行数相对于总量来说极少高筛选率。ClickHouse的列式存储和向量化执行引擎在处理这类带有强力过滤条件的聚合查询时性能是碾压式的。它不像HBase那样擅长单行随机读写也不像ES那样在全文检索上独占鳌头但在我们这种“海量数据快速过滤”的场景下它是最优解。其次成本与效率的平衡。自建ClickHouse集群的硬件成本尤其是SSD和运维复杂度相对于其带来的性能提升ROI投资回报率非常高。它的数据压缩比惊人通常能达到10:1甚至更高这意味着存储万亿条记录所需的物理磁盘空间远小于其他方案。同时其单机性能强悍在合理分片和副本设计下可以通过增加节点近乎线性地提升查询能力。注意ClickHouse并非银弹。它对高频、小事务的更新/删除操作支持很差虽然现在有DELETE和UPDATE但代价高昂。我们的轨迹数据一旦入库基本以追加为主极少修改完美避开了它的短板。2.2 空间索引的抉择从Geohash到S2轨迹数据检索的核心是空间查询。如何高效判断一个点是否在一个多边形区域内如何快速找到彼此邻近的点这需要空间索引。早期我们尝试过Geohash。它将二维的经纬度编码成一维字符串前缀匹配可以快速找到大致在同一区域的数据。但它有个致命问题边界效应。两个地理上非常接近的点如果恰好位于Geohash网格的边界两侧它们的Geohash编码可能完全不同导致查询遗漏。这对于精度要求高的地理围栏查询是不可接受的。因此我们转向了Google开源的S2 Geometry Library。S2将地球表面投影到一个立方体上再细分为层次化的Cell。每个Cell有一个唯一的64位Cell ID。相比于GeohashS2具有多项优势形状更规整S2的Cell在真实地球表面的形状更接近方形变形较小。层次化覆盖从Level 0整个地球到Level 30约1cm²可以自由选择精度。对于车辆轨迹Level 15边长约300米或Level 16边长约150米通常是不错的折中选择。强大的几何关系计算S2原生提供了判断点、线、面之间包含、相交等关系的功能算法高效且准确。我们的策略是在数据入库时实时计算每个轨迹点经纬度对应的S2 Cell ID比如Level 16并将这个Cell ID作为一列存储在ClickHouse中。当查询一个多边形区域时我们先利用S2库计算出覆盖这个多边形所有可能相交的Cell ID集合这是一个范围然后在ClickHouse中利用WHERE s2_cell_id IN (...)或对Cell ID范围进行WHERE s2_cell_id BETWEEN ... AND ...来快速筛选出候选点集最后再用S2的精确几何计算进行二次过滤。这种“二级过滤”机制将耗时的几何计算从全表扫描中解放出来先通过索引快速缩小数据范围是达到秒级响应的关键。2.3 数据模型设计兼顾查询效率与存储在ClickHouse中表结构设计直接决定了性能上限。我们的核心轨迹表trajectory_fact大致设计如下CREATE TABLE trajectory_fact_local ( vehicle_id String, -- 车辆唯一标识 timestamp DateTime64(3, Asia/Shanghai), -- 精确到毫秒的时间戳 longitude Float64, -- 经度 latitude Float64, -- 纬度 s2_cell_id UInt64, -- S2 Cell ID (Level 16) date Date MATERIALIZED toDate(timestamp), -- 用于分区的衍生列 hour UInt8 MATERIALIZED toHour(timestamp) -- 用于进一步分区的衍生列 ) ENGINE MergeTree() PARTITION BY (date, hour) ORDER BY (s2_cell_id, timestamp, vehicle_id) SETTINGS index_granularity 8192;设计要点解析分区键PARTITION BY我们按date和hour进行分区。这是时间序列数据的经典做法。每天、每小时的数据独立存储当查询指定时间范围时ClickHouse可以快速分区裁剪只加载相关分区文件极大减少IO。例如查询“2023-10-01 10点到11点”的数据只会扫描2023-10-01_10这一个分区。排序键ORDER BY这是ClickHouse性能的灵魂。我们将其设置为(s2_cell_id, timestamp, vehicle_id)。s2_cell_id在最前面因为我们的核心查询总是基于空间范围。将S2 Cell ID作为主排序键可以保证相同或相邻Cell的数据在物理磁盘上连续存储。当执行WHERE s2_cell_id IN (...)查询时ClickHouse可以进行高效的数据跳跃扫描直接定位到相关的数据块避免全表扫描。timestamp紧随其后在空间筛选后时间过滤是第二高频操作。这样的排序使得在指定了S2 Cell ID后按时间范围查找也非常高效。vehicle_id用于优化按车辆ID查询的场景。索引粒度index_granularity默认为8192。这意味着每8192行数据形成一个“数据块”并为之创建一个稀疏索引条目记录每个数据块中排序键的最小最大值。较小的粒度如1024会创建更密集的索引查询时定位更精准但索引文件会变大消耗更多内存。经过测试对于我们的数据规模和查询模式8192是一个平衡点。3. 核心检索流程的深度拆解与实现有了底层存储和索引接下来就是实现具体的检索流程。一个完整的“多边形区域轨迹查询”流程远不止一句SQL那么简单。3.1 查询预处理从地理围栏到S2 Cell列表业务方传来的通常是一个GeoJSON格式的多边形坐标串。第一步是将其转化为ClickHouse能高效利用的S2 Cell ID列表。import s2sphere as s2 def polygon_to_s2_cell_ids(polygon_coords, s2_level16): 将多边形坐标转换为覆盖该多边形的S2 Cell ID列表。 polygon_coords: 列表格式如 [[lng1, lat1], [lng2, lat2], ...] s2_level: S2 Cell的级别 # 1. 创建S2多边形对象 points [s2.LatLng.from_degrees(lat, lng) for lng, lat in polygon_coords] s2_polygon s2.S2Polygon(s2.S2Loop(points)) # 2. 获取覆盖多边形的S2 Cell覆盖器 coverer s2.S2RegionCoverer() coverer.min_level s2_level coverer.max_level s2_level coverer.max_cells 100 # 控制返回的Cell数量防止过多 # 3. 获取Cell ID列表 covering coverer.get_covering(s2_polygon) cell_ids [cell.id() for cell in covering] return cell_ids这个过程有几个关键参数和陷阱s2_level的选择级别越高单个Cell面积越小覆盖同一区域需要的Cell数量越多查询时的索引过滤精度越高但IN列表也会更长可能影响查询性能。需要根据业务对精度的要求和查询延迟进行权衡。我们通过A/B测试发现对于城市级围栏查询Level 16在精度和性能上综合表现最好。max_cells的限制必须设置一个上限防止一个巨大的多边形如跨省产生数万个Cell ID导致查询条件过于庞大甚至失败。如果覆盖的Cell数量超过上限S2库会返回一个近似覆盖可能用更高级别的大Cell来覆盖。这时需要在日志中告警提示业务方围栏可能过大或需要调整参数。3.2 ClickHouse查询构建与执行拿到S2 Cell ID列表后就可以构建ClickHouse查询了。查询的核心思路是“二级过滤”。-- 假设我们得到的S2 Cell ID列表是 [9926595695168651264, 9926595743884648448, ...] WITH target_cell_ids AS (SELECT arrayJoin([9926595695168651264, 9926595743884648448]) AS id) SELECT vehicle_id, timestamp, longitude, latitude FROM trajectory_fact WHERE date 2023-10-01 -- 利用分区键进行第一层快速裁剪 AND hour 10 -- 利用分区键进行第一层快速裁剪 AND s2_cell_id IN target_cell_ids -- 利用排序键进行第二层快速数据跳跃 -- 第三层精确的空间几何判断在ClickHouse中调用S2函数或使用预先计算好的关系 AND s2Contains({{polygon_wkt}}, S2Point(longitude, latitude)) 1 ORDER BY timestamp LIMIT 1000查询优化点详解使用WITH子句和数组将Cell ID列表定义为一个公共表表达式(CTE)或直接使用数组可以使查询语句更清晰且ClickHouse对数组的IN查询有优化。分区键优先务必把对分区键date和hour的过滤条件放在最前面确保分区裁剪最先发生。函数计算下推s2Contains是精确判断点是否在多边形内的函数。我们将其放在WHERE条件中ClickHouse会在读取数据行时进行计算过滤。但要注意这个计算是CPU密集型的所以前两层过滤分区S2索引必须足够高效将需要精确计算的数据量降到万级甚至千级以下。**避免SELECT ***只查询需要的列。ClickHouse是列式存储读取不需要的列会造成不必要的IO浪费。3.3 应对超大规模结果集分页与流式返回当查询区域很大或时间跨度很长时命中的轨迹点可能多达数百万甚至上千万。一次性返回所有结果不仅耗时长而且可能压垮客户端和网络。我们实现了两种策略策略一基于时间戳的“游标”分页对于需要深度翻页的场景我们不让用户传page和size而是传last_timestamp和last_vehicle_id。WHERE ... AND (timestamp 2023-10-01 10:00:00 OR (timestamp 2023-10-01 10:00:00 AND vehicle_id last_vid)) ORDER BY timestamp, vehicle_id LIMIT 1000这种方式利用排序键进行过滤每次查询都能高效地利用索引找到“下一页”的起点性能稳定不受翻页深度影响。策略二流式HTTP响应与聚合对于需要实时在地图上渲染大量轨迹点的前端应用我们采用HTTP Chunked Encoding进行流式返回。ClickHouse支持FORMAT JSONEachRow格式配合settings stream_like_format1可以边查询边输出。后端服务接收到一行就转发一行给前端前端进行增量渲染用户体验是数据逐渐加载出来而不是长时间等待后一次性爆炸式呈现。4. 性能调优与稳定性保障实战系统上线后真正的挑战才开始。随着数据量从十亿到千亿再到万亿各种性能瓶颈和稳定性问题接踵而至。4.1 常见性能瓶颈分析与解决瓶颈一S2 Cell ID列表过长导致查询解析慢现象查询一个形状不规则的大型围栏产生的S2 Cell ID多达上万个生成的SQL中IN子句极其庞大ClickHouse解析SQL的时间甚至超过执行时间。解决方案合并相邻Cell对生成的Cell ID列表进行预处理将连续的、可以合并成更大级别Cell的ID进行合并减少列表长度。例如四个连续的Level 16 Cell可以合并为一个Level 15 Cell。使用临时表将Cell ID列表先插入到一张ClickHouse的临时内存表ENGINE Memory中然后使用JOIN来代替IN查询。ClickHouse对JOIN的优化有时比超长列表的IN更好。业务侧优化与产品经理沟通为地理围栏设置最大面积限制或引导用户绘制更规整的查询区域。瓶颈二热点数据分区查询压力大现象业务总是查询最近一天的数据导致最新的一两个分区承受了绝大部分的查询流量该分区所在磁盘IOPS和CPU使用率持续高位。解决方案冷热数据分层将超过30天的历史数据转移到存储成本更低、查询频率较低的机械硬盘卷上甚至转移到对象存储如S3并通过ClickHouse的S3磁盘类型或外部表来访问。最新数据放在高性能SSD上。加强监控对每个分区的查询QPS、扫描行数进行监控及时发现热点。考虑更细粒度分区如果单小时数据量仍然巨大如超过1亿条可以考虑按15分钟甚至更细粒度分区让查询压力更分散但需注意分区总数不宜过多通常建议不超过1万。瓶颈三并发查询下的资源争抢现象多个复杂空间查询同时执行导致ClickHouse节点内存max_memory_usage爆掉查询被终止Memory limit exceeded。解决方案设置用户资源配额在ClickHouse的users.xml中为不同的业务用户组设置不同的max_memory_usage、max_execution_time等限制防止单个劣质查询拖垮整个集群。查询队列与熔断在应用层实现查询队列控制发往ClickHouse的并发数。对于耗时超过一定阈值的查询自动熔断并返回降级结果如提示“查询超时请缩小范围”。优化查询本身这是根本。通过EXPLAIN语句分析查询计划确保索引被正确使用。避免在WHERE条件中对列进行函数计算如WHERE formatDateTime(timestamp, %Y-%m-%d) ...这会导致索引失效。4.2 数据写入与更新的权衡轨迹数据是高速写入的。我们面临的是每秒数十万甚至上百万点的写入压力。写入优化批量写入绝对禁止单条INSERT。我们使用异步缓冲队列积攒一定数量如1万条或达到一定时间窗口如1秒后批量写入ClickHouse。ClickHouse的MergeTree引擎对批量写入友好能大幅减少parts数量减轻后台合并压力。使用合适的表引擎对于实时写入我们使用了MergeTree的变种ReplicatedReplacingMergeTree在保证数据最终一致性的同时通过副本提高了可用性。写入时直接写入本地表和所有副本的Distributed表逻辑过于沉重我们采用了“先写本地再通过ZooKeeper异步复制”的模式。监控写入延迟与堆积建立完善的监控关注parts数量、merges速度以及buffer表的堆积情况。一旦发现写入延迟增大或parts数量暴涨需要立即排查是数据源问题、网络问题还是ClickHouse集群负载问题。关于数据更新与删除轨迹数据原则上不更新。但业务存在纠偏需求如GPS漂移修正。我们采用“标记删除重新插入”的方式。新增一个is_deletedUInt8字段默认0。需要删除或更新时先执行ALTER TABLE ... UPDATE is_deleted 1 WHERE ...这是一个异步、耗资源的操作需在低峰期进行。将修正后的新数据直接插入。在查询时始终加上WHERE is_deleted 0的条件。定期如每月使用ALTER TABLE ... DELETE WHERE is_deleted 1进行物理清理或者通过创建新分区表并迁移有效数据的方式来“重整”数据。5. 监控、告警与故障排查体系一个稳定运行的系统离不开可观测性。我们建立了一套从基础设施到业务指标的立体监控体系。5.1 核心监控指标监控类别具体指标告警阈值说明集群健康ZooKeeper连接状态、副本同步延迟任何异常基础影响分布式表写入和复制。节点资源CPU使用率、内存使用率、磁盘使用率/IOPS、网络带宽85%持续5分钟资源瓶颈的直接体现。查询性能查询平均耗时(P99)、查询QPS、慢查询数量P99 2s, 慢查询10个/分钟业务体验的生命线。写入性能每秒写入行数、写入耗时、Buffer表堆积行数写入耗时1s, Buffer堆积100万影响数据实时性。Merge状态待合并parts数量、后台合并线程状态parts数500 合并线程异常parts过多影响查询性能。业务指标核心接口响应时间、错误率响应时间1s错误率0.1%从用户视角监控。5.2 典型故障排查实录案例一查询突然变慢但CPU/内存不高现象业务反馈某个常用围栏查询从200ms飙升到10s以上。登录服务器查看ClickHouse节点资源使用率正常。排查步骤登录ClickHouse客户端执行SHOW PROCESSLIST查看当前正在运行的查询。发现有几个陌生的复杂聚合查询长时间运行。使用EXPLAIN分析慢查询发现其扫描了全表数亿行数据s2_cell_id索引未生效。原因是查询条件中使用了OR连接了多个不相关的围栏且写法导致索引优化器失效。根因业务方新上线了一个功能并发起了多个跨区域的不合理查询。解决短期通过KILL QUERY终止问题查询。长期在查询网关层对SQL进行简单的语法分析拦截明显不合理的查询如没有有效分区键或索引键过滤并推动业务方修改查询逻辑将大查询拆分为多个可以利用索引的小查询并行执行。案例二磁盘空间告警但数据量预期内现象监控显示某节点磁盘使用率一周内从60%快速增长到95%。排查步骤使用SELECT table, sum(bytes_on_disk) FROM system.parts GROUP BY table ORDER BY sum(bytes_on_disk) DESC查看各表磁盘占用。发现某张日志表异常巨大。检查该表的TTL生存时间设置发现由于配置错误TTL并未生效历史数据从未被删除。检查该表的parts数量SELECT COUNT(*) FROM system.parts WHERE table log_table发现数量高达数万远超正常水平通常几百以内。说明大量小批量写入产生了大量parts且合并可能受阻。根因TTL配置错误 写入过于零散。解决首先修正TTL配置。其次因为直接删除大量数据可能引发长时间锁表我们选择创建一个新的分区表将需要保留的数据插入新表然后原子性地RENAME交换表名。最后优化该数据源的写入逻辑增加批量写入的缓冲大小减少parts产生。经过这些实战的打磨我们的图灵平台轨迹检索服务终于能够稳定地应对万亿级数据下的秒级查询挑战。这个过程没有一劳永逸的“黑科技”更多的是对底层原理的深刻理解、严谨的架构设计、精细的工程实现以及持续不断的监控优化。每一个参数的选择每一次查询的编写都需要结合具体的数据分布和业务场景反复权衡。希望我们踩过的这些坑和总结出的经验能为同样面临海量时空数据检索难题的团队提供一些切实可行的参考。
返回列表