You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

使用IN操作符查询Cassandra时遇NoHostAvailable/ConnectionBusy问题求助

Cassandra查询出现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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.21 01:37:19