Skip to content

一、Cassandra 与时序大数据的契约 ​

在工业互联网、监控遥测或物联网(IoT)场景中,每秒可能产生数万甚至数百万个测点的振动、温度与指标数据。传统关系型数据库(如 MySQL、PG)面对如此高频的写入会迅速遇到磁盘 I/O 和索引维护瓶颈。

Apache Cassandra 凭借其基于 LSM-Tree(Log-Structured Merge-tree)追加写入的高吞吐特性与去中心化的环形哈希(Consistent Hashing)架构,成为海量时序存储的核心基础设施。

然而,很多习惯了关系型 SQL 语法的工程师在初接触 CQL(Cassandra Query Language)时,往往会踩入严重的性能误区。


二、核心存储模型:分区键 vs 聚簇列 ​

Cassandra 的主键(PRIMARY KEY)设计决定了数据在物理集群中的分布拓扑:

cql
CREATE TABLE device_timeseries (
    device_id uuid,
    metric_code text,
    event_time timestamp,
    metric_value double,
    PRIMARY KEY ((device_id, metric_code), event_time)
) WITH CLUSTERING ORDER BY (event_time DESC);
  1. 复合分区键(Partition Key):((device_id, metric_code))
    • 作用:哈希散列计算,决定数据归属于 Cassandra 集群的哪一个虚拟节点(vNode)和物理机器上。
  2. 聚簇列(Clustering Column):event_time
    • 作用:在同一台机器的单个分区(SSTable)内部,数据物理上按照时间倒序(DESC)顺序紧凑排列。

三、为什么严禁在线滥用 ALLOW FILTERING? ​

在执行 CQL 查询时,如果 WHERE 条件没有命中完整的分区键,执行时会触发报错:

cql
-- 报错:Cannot execute this query to find data without knowing complete partition key
SELECT * FROM device_timeseries WHERE metric_value > 80.0;

很多新手看到提示,直接在末尾加上 ALLOW FILTERING:

cql
-- 危险用法:强行开启全表扫描
SELECT * FROM device_timeseries WHERE metric_value > 80.0 ALLOW FILTERING;

致命危害分析: ​

  • 触发跨节点全量扫表:协调节点(Coordinator)必须向集群中的所有物理节点发送读请求,每个节点都会遍历本地所有的 SSTable,产生海量的磁盘随机读与网络广播;
  • GC 停顿与节点宕机:全表扫描会将数以百万计的墓碑(Tombstones)与无用行载入 JVM 堆内存,极易引发长时间 JVM Full GC 甚至触发 OOM 崩溃;
  • 正确原则:Cassandra 是为查询而设计表模型(Query-driven Modeling)。如果业务需要按 metric_value 查询,应建立物化视图、专用倒排二级索引或辅助查询表,而不是靠 ALLOW FILTERING 强行跨分区扫表。

四、Python 驱动连接与参数化查询 ​

使用官方推荐的 cassandra-driver 封装数据库会话:

python
from cassandra.cluster import Cluster, ExecutionProfile, EXEC_PROFILE_DEFAULT
from cassandra.policies import DCAwareRoundRobinPolicy, RetryPolicy
from cassandra.query import SimpleStatement
from cassandra import ConsistencyLevel

class CassandraManager:
    _cluster = None
    _session = None

    @classmethod
    def get_session(cls, contact_points=["127.0.0.1"], port=9042, keyspace="iot_metrics"):
        if cls._session is None:
            # 配置局部感知轮询策略与一致性级别
            profile = ExecutionProfile(
                load_balancing_policy=DCAwareRoundRobinPolicy(local_dc="datacenter1"),
                retry_policy=RetryPolicy(),
                consistency_level=ConsistencyLevel.LOCAL_QUORUM # 兼顾强一致与低时延
            )
            cls._cluster = Cluster(contact_points=contact_points, port=port,
                                   execution_profiles={EXEC_PROFILE_DEFAULT: profile})
            cls._session = cls._cluster.connect(keyspace)
        return cls._session

预编译查询(PreparedStatement)实操 ​

python
def fetch_recent_metrics(session, device_id, metric_code, limit=100):
    # 必须对常用查询进行预编译,Coordinator 会缓存查询分析树
    cql = """
        SELECT event_time, metric_value
        FROM device_timeseries
        WHERE device_id = ? AND metric_code = ?
        LIMIT ?;
    """
    prepared = session.prepare(cql)
    rows = session.execute(prepared, (device_id, metric_code, limit))
    
    results = []
    for row in rows:
        results.append({
            "time": row.event_time.strftime("%Y-%m-%d %H:%M:%S"),
            "value": row.metric_value
        })
    return results

测试开发工程师 · 专注自动化与系统架构 | 邮箱: hansblog@atumsoul.win