ScyllaDB百万级数据聚合查询超时的配置优化咨询
我创建了如下表结构:
create table if not exists keyspace_test.table_test ( id int, date text, val float, primary key (id, date) ) with caching = {'keys': 'ALL', 'rows_per_partition': 'ALL'} and compaction = {'class': 'SizeTieredCompactionStrategy'} and compression = {'sstable_compression': 'org.apache.cassandra.io.compress.LZ4Compressor'} and dclocal_read_repair_chance = 0 and speculative_retry = '99.0PERCENTILE' and read_repair_chance = 1;
导入1200万行数据后,执行以下聚合查询统计行数并求和:
SELECT COUNT(*), SUM(val) FROM keyspace_test.table_test
出现错误:
Cassandra timeout during read query at consistency ONE (1 responses were required but only 0 replica responded)
添加USING TIMEOUT 180s;后仍然报错:
Timed out waiting for server response
集群配置为两个数据中心,各4台服务器:
# docker exec -it scylla-120 nodetool status Datacenter: dc2 =============== Status=Up/Down |/ State=Normal/Leaving/Joining/Moving -- Address Load Tokens Owns Host ID Rack UN 10.3.192.25 79.04 GB 256 ? 5975a143fec6 Rack1 UN 10.3.192.24 74.2 GB 256 ? 61dc1cfd3e92 Rack1 UN 10.3.192.22 88.21 GB 256 ? 0d24d52d6b0a Rack1 UN 10.3.192.23 63.41 GB 256 ? 962s266518ee Rack1 Datacenter: dc3 =============== Status=Up/Down |/ State=Normal/Leaving/Joining/Moving -- Address Load Tokens Owns Host ID Rack UN 34.77.78.21 83.5 GB 256 ? 5112f248dd38 Rack1 UN 34.77.78.20 59.87 GB 256 ? e8db897ca33b Rack1 UN 34.77.78.48 81.32 GB 256 ? cb88bd9326db Rack1 UN 34.77.78.47 79.8 GB 256 ? 562a721d4b77 Rack1 Note: Non-system keyspaces don't have the same replication settings, effective ownership information is meaningless
Keyspace创建语句:
CREATE KEYSPACE keyspace_test WITH replication = { 'class' : 'NetworkTopologyStrategy', 'dc2' : 3, 'dc3' : 3};
请问如何配置ScyllaDB以处理百万级数据的这类聚合查询?
1. 避免直接全表扫描聚合
ScyllaDB是分布式列存数据库,核心设计是支持低延迟的点查询和范围查询,全表扫描的COUNT(*)/SUM()这类聚合查询本质上需要遍历所有节点的所有数据,天生效率极低,容易超时,这不是它的优势场景。
2. 预计算聚合值(推荐方案)
方式一:使用物化视图
创建专门存储聚合结果的物化视图,让ScyllaDB自动维护聚合值:
CREATE MATERIALIZED VIEW keyspace_test.table_test_agg AS SELECT id, COUNT(*) as total_rows, SUM(val) as total_val FROM keyspace_test.table_test GROUP BY id;
如果需要全局聚合(不按id分组),可以设置固定分组键,比如:
CREATE TABLE keyspace_test.global_agg ( agg_key text PRIMARY KEY, total_rows bigint, total_val float );
然后通过应用程序或定时任务(比如ScyllaDB Task Manager)定期计算并更新这个表的聚合值,避免每次查询都全表扫描。
方式二:ETL定时计算
用Spark、Dask等工具定期读取数据,计算聚合结果后写入专门的统计表,业务查询直接读取该统计表即可。
3. 调整查询参数和一致性级别
如果必须执行即时全表聚合,可以做以下调整:
- 设置LOCAL_ONE一致性级别:避免跨数据中心查询,只从本地DC获取数据,减少网络延迟:
SELECT COUNT(*), SUM(val) FROM keyspace_test.table_test USING CONSISTENCY LOCAL_ONE AND TIMEOUT 300s;
- 路由到单个节点:使用
TOKEN函数让查询只路由到一个节点,由该节点协调所有分片的聚合(注意:会给该节点带来较大压力):
SELECT COUNT(*), SUM(val) FROM keyspace_test.table_test WHERE TOKEN(id) >= 0 AND TOKEN(id) < 2^63 USING CONSISTENCY LOCAL_ONE;
4. 优化表和集群配置
- 调整compaction策略:对于需要频繁扫描的表,改用
LeveledCompactionStrategy,减少扫描时需要读取的SSTable数量,提升扫描效率:
ALTER TABLE keyspace_test.table_test WITH compaction = {'class': 'LeveledCompactionStrategy'};
- 关闭read_repair_chance:当前表设置
read_repair_chance = 1,会导致每次查询都触发读修复,极大增加查询开销,建议改为0:
ALTER TABLE keyspace_test.table_test WITH read_repair_chance = 0;
- 调整缓存设置:当前
rows_per_partition = 'ALL'会缓存整个分区的数据,大分区可能占用过多内存,建议改为固定行数或关闭:
ALTER TABLE keyspace_test.table_test WITH caching = {'keys': 'ALL', 'rows_per_partition': '100'};
- 检查节点负载:从
nodetool status看节点负载在60-90GB之间,确保节点有足够的内存和CPU资源,避免资源瓶颈导致超时。
5. 使用ScyllaDB的分析扩展
如果需要频繁做这类分析查询,可以部署ScyllaDB Analytics,它集成了Spark,能高效处理分布式数据的聚合分析,适合大数据量的统计场景。
内容的提问来源于stack exchange,提问作者Pamungkas Jayuda

