使用IN操作符查询Cassandra时遇NoHostAvailable/ConnectionBusy问题求助
问题背景
使用带IN操作符的查询从Cassandra获取特定时间段数据,查询语句如下:
"SELECT * FROM "+ DATABASES['default']['NAME'] +".perform_stats WHERE device_id IN (" + ','.join(device_id_lst) + " AND created_at <= " + created_at +" and created_at >= " start_time + ";"
随机出现以下错误:
cassandra.cluster.NoHostAvailable: ('Unable to complete the operation against any hosts', {<Host: host:9042 datacenter1>: ConnectionBusy('Connection host:9042 is overloaded')}) File "/opt/app-root/lib64/python3.8/site-packages/cassandra/cqlengine/query.py", line 515, in __iter__ File "/opt/app-root/lib64/python3.8/site-packages/cassandra/cqlengine/query.py", line 472, in _execute_query File "/opt/app-root/lib64/python3.8/site-packages/cassandra/cqlengine/query.py", line 404, in _execute File "/opt/app-root/lib64/python3.8/site-packages/cassandra/cqlengine/query.py", line 1531, in _execute_statement File "/opt/app-root/lib64/python3.8/site-packages/cassandra/cqlengine/connection.py", line 345, in execute File "cassandra/cluster.py", line 2618, in cassandra.cluster.Session.execute File "cassandra/cluster.py", line 4877, in cassandra.cluster.ResponseFuture.result cassandra.cluster.NoHostAvailable: ('Unable to complete the operation against any hosts', {<Host: host:9042 datacenter1>: ConnectionBusy('Connection host:9042 is overloaded')})
单条查询针对10个设备15分钟内的数据(每个设备约4条记录,返回约40条结果),但系统内存缓存急剧上升(堆内存占用6/8GB),且错误随机触发。已尝试增加内存、减少单查询设备数至10、替换为原生查询,均无效果。
表结构
CREATE TABLE sample.perform_stats ( device_id uuid, created_at bigint, stats_data text, PRIMARY KEY (device_id, created_at) ) WITH CLUSTERING ORDER BY (created_at DESC) AND bloom_filter_fp_chance = 0.01 AND caching = {'keys': 'ALL', 'rows_per_partition': 'NONE'} AND comment = '' AND compaction = {'class': 'org.apache.cassandra.db.compaction.SizeTieredCompactionStrategy', 'max_threshold': '32', 'min_threshold': '4'} AND compression = {'chunk_length_in_kb': '64', 'class': 'org.apache.cassandra.io.compress.LZ4Compressor'} AND crc_check_chance = 1.0 AND dclocal_read_repair_chance = 0.1 AND default_time_to_live = 0 AND gc_grace_seconds = 864000 AND max_index_interval = 2048 AND memtable_flush_period_in_ms = 0 AND min_index_interval = 128 AND read_repair_chance = 0.0 AND speculative_retry = '99PERCENTILE';
环境信息
- CentOS Version: 7.9
- Cassandra version: 4.1.0
- Java Version: 11.0.18
- Python Version: 3.6.8
- Django: 3.1.4
- django-cassandra-engine: 1.6.1
- cassandra-driver==3.24.0
可能原因与解决方法
1. 修复查询语句语法错误
原查询存在语法问题(IN括号未闭合、字符串拼接缺失加号),先修正语句避免执行异常:
# 使用f-string简化拼接,避免语法错误 query = f"SELECT * FROM {DATABASES['default']['NAME']}.perform_stats WHERE device_id IN ({','.join(device_id_lst)}) AND created_at <= {created_at} AND created_at >= {start_time};"
2. 调整客户端连接池配置
ConnectionBusy错误通常与客户端连接池耗尽有关,cassandra-driver默认连接数较低,需调高节点连接上限:
原生驱动配置
from cassandra.cluster import Cluster cluster = Cluster( contact_points=['host'], port=9042, protocol_version=4, connections_per_host=10, # 默认2,适当调高 max_connections_per_host=20, pool_timeout=30 # 延长超时,避免快速触发连接繁忙 ) session = cluster.connect()
Django Cassandra Engine配置
在settings.py中添加连接池参数:
DATABASES = { 'default': { 'ENGINE': 'django_cassandra_engine', 'NAME': 'sample', 'HOST': 'host', 'OPTIONS': { 'connection': { 'connections_per_host': 10, 'max_connections_per_host': 20, 'pool_timeout': 30, } } } }
3. 替换IN查询为并行单分区查询
Cassandra的IN查询会将请求分发到多个节点(每个device_id是独立分区),即使仅10个值也可能触发跨节点连接压力。建议改为并行查询单个分区:
from concurrent.futures import ThreadPoolExecutor from cassandra.query import SimpleStatement def fetch_single_device(session, device_id, created_at, start_time): query = f"SELECT * FROM sample.perform_stats WHERE device_id = {device_id} AND created_at <= {created_at} AND created_at >= {start_time};" # 设置fetch_size分批拉取,减少内存占用 stmt = SimpleStatement(query, fetch_size=10) return list(session.execute(stmt)) # 并行执行查询,控制并发数避免客户端过载 with ThreadPoolExecutor(max_workers=10) as executor: futures = [ executor.submit(fetch_single_device, session, did, created_at, start_time) for did in device_id_lst ] # 收集结果 all_results = [f.result() for f in futures]
这种方式更贴合Cassandra分布式架构,分散节点负载。
4. 优化Cassandra节点与表配置
调整表缓存策略
当前表配置rows_per_partition: 'NONE',无法缓存常用行,导致频繁磁盘IO。针对时间序列查询,设置缓存最近的分区行:
ALTER TABLE sample.perform_stats WITH caching = {'keys': 'ALL', 'rows_per_partition': '100'};
切换压缩策略
当前使用的SizeTieredCompactionStrategy(STCS)不适合时间序列数据,建议改用TimeWindowCompactionStrategy(TWCS),按时间窗口合并SSTable,降低compaction压力:
ALTER TABLE sample.perform_stats WITH compaction = { 'class': 'org.apache.cassandra.db.compaction.TimeWindowCompactionStrategy', 'compaction_window_size': 15, 'compaction_window_unit': 'MINUTES' };
优化堆内存与GC配置
检查cassandra-env.sh中的堆内存设置,Cassandra 4.1建议堆内存为物理内存的1/4~1/2,确保MAX_HEAP_SIZE和HEAP_NEWSIZE合理,避免频繁Full GC。同时使用nodetool gcstats监控GC情况,调整JVM参数。
5. 客户端内存优化
- 分批拉取结果:使用
fetch_size参数,避免一次性加载所有结果到内存:
stmt = SimpleStatement(query, fetch_size=10) for row in session.execute(stmt): # 逐行处理,减少内存占用 process_row(row)
- 减少不必要字段:在Django ORM中使用
only()或values()仅获取需要的字段,降低对象内存开销:
# 仅加载所需字段 results = PerformStats.objects.filter( device_id__in=device_id_lst, created_at__range=(start_time, created_at) ).only('device_id', 'created_at', 'stats_data')
6. 监控与排查工具
- 使用
nodetool status查看节点状态,nodetool tpstats检查线程池是否有请求堆积; - 执行
TRACING ON后运行查询,跟踪请求的节点通信路径,定位瓶颈; - 使用cassandra-driver的metrics功能监控连接池使用情况:
from cassandra import metrics print(metrics.get_metrics())
内容的提问来源于stack exchange,提问作者Deepak Keshari

