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

Cassandra Python驱动无法返回全量行数COUNT()的问题排查求助

Cassandra插入14000行仅返回1500行问题排查与修复

问题背景

使用DataStax免费版Cassandra(25MB存储限制)插入约14000行DataFrame(数据集约1.5MB),插入、查询操作均无报错,但查询仅返回约1500行。需定位故障环节,解决全量插入/查询问题。


插入代码(cassandraDBLoad.py)

def progressbar(it, prefix="", size=60, out=sys.stdout): # Python3.3+
    count = len(it)
    def show(j):
        x = int(size*j/count)
        print("{}[{}{}] {}/{}".format(prefix, u"█"*x, "."*(size-x), j, count), 
                end='\r', file=out, flush=True)
    show(0)
    for i, item in enumerate(it):
        yield item
        show(i+1)
    print("\n", flush=True, file=out)

def cassandraDBLoad(config_path):
    try:
        config = read_params(config_path)

        execution_profile = ExecutionProfile(request_timeout=10)
        cassandra_config = {'secure_connect_bundle': config["connect_cassandra"]["cloud_config"]}
        auth_provider = PlainTextAuthProvider(
                config["connect_cassandra"]["client_id"],
                config["connect_cassandra"]["client_secret"]
                )
        cluster = Cluster(cloud=cassandra_config, auth_provider=auth_provider)
        session = cluster.connect()
        session.default_timeout = None
        connect_db = session.execute("select release_version from system.local")
        set_keyspace = session.set_keyspace(config["cassandra_db"]["keyspace"])
        
        table_ = config["cassandra_db"]["data_table"]
        define_columns = config["cassandra_db"]["define_columns"]
        
        create_table = f"CREATE TABLE IF NOT EXISTS {table_}({define_columns});"
        start_create = time.process_time()
        table_result = session.execute(create_table)
        
        train = pd.read_csv(config["data_source"]["train_source"])
        test = pd.read_csv(config["data_source"]["test_source"])
    
        #Combine test and train into one file
        train['source']='train'
        test['source']='test'
        df = pd.concat([train, test],ignore_index=True)
        df = df.fillna('NA')
        columns = list(df)
        for col in columns:
            df[col] = df[col].map(str)
        
        columns = config["cassandra_db"]["columns"]
        insert_qry = f"INSERT INTO {table_}({columns}) VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?) IF NOT EXISTS"
        statement = session.prepare(insert_qry)
        
        start_insert = time.process_time()
        batch = BatchStatement()
        for i in progressbar(range(len(df)), "Inserting: ", 40):
            time.sleep(0.1)            
            session.execute_async(
                statement,
                    [
                        df.iat[i,0], df.iat[i,1], df.iat[i,2], df.iat[i,3], df.iat[i,4], df.iat[i,5], 
                        df.iat[i,6], df.iat[i,7], df.iat[i,8], df.iat[i,9], df.iat[i,10], df.iat[i,11], 
                        df.iat[i,12]
                    ]
                )
            print("Time taken to insert df - " + str((time.process_time() - start_insert)/60) + " minutes")

    except Exception as e:
        raise Exception("(cassandraDBLoad): Something went wrong in the CassandraDB Load operations\n" + str(e))

插入耗时约30分钟,机器配置:12GB内存、2核CPU


查询代码(preprocess_data.py)

def pandas_factory(colnames, rows):
    return pd.DataFrame(rows, columns=colnames)

def preprocess_data(config_path):
    try:
        config = read_params(config_path)
        cassandra_config = {'secure_connect_bundle': config["connect_cassandra"]["cloud_config"]}
        auth_provider = PlainTextAuthProvider(
                config["connect_cassandra"]["client_id"],
                config["connect_cassandra"]["client_secret"]
                )
        cluster = Cluster(cloud=cassandra_config, auth_provider=auth_provider)
        session = cluster.connect()
        session.set_keyspace(config["cassandra_db"]["keyspace"])
        session.row_factory = pandas_factory
        #session.default_fetch_size = None

        count_query = f"SELECT COUNT(*) from {config['cassandra_db']['data_table']} LIMIT 20000"
        count_rslt = session.execute(count_query, timeout=None)
        print(count_rslt._current_rows)
        query = f"SELECT * from {config['cassandra_db']['data_table']}"
        simple_statement = SimpleStatement(query, consistency_level=ConsistencyLevel.ONE, fetch_size=None)
        execute_result = session.execute(simple_statement, timeout=None)
        data = execute_result._current_rows
  
        print(data.shape)
        
    except Exception as e:
        raise Exception("(preprocessData): " + str(e))

故障定位与修复方案

1. 插入代码核心问题

  • 异步操作未处理结果:execute_async发起请求后未等待结果,插入失败不会抛出异常,导致大量数据丢失。
  • 无效sleep拖慢速度:time.sleep(0.1)导致光等待就耗时23分钟,且无任何必要。
  • IF NOT EXISTS导致主键冲突丢弃:如果数据集存在重复主键,这部分行不会被插入,需先确认主键唯一性。
  • 未使用批量插入:定义了BatchStatement但未实际使用,单条插入效率极低,易触发Cassandra限流。

插入代码修复示例

start_insert = time.process_time()
futures = []
# 移除sleep,收集所有异步请求结果
for i in progressbar(range(len(df)), "Inserting: ", 40):
    future = session.execute_async(
        statement,
        [
            df.iat[i,0], df.iat[i,1], df.iat[i,2], df.iat[i,3], df.iat[i,4], df.iat[i,5], 
            df.iat[i,6], df.iat[i,7], df.iat[i,8], df.iat[i,9], df.iat[i,10], df.iat[i,11], 
            df.iat[i,12]
        ]
    )
    futures.append(future)

# 等待所有请求完成,捕获失败请求
for idx, future in enumerate(futures):
    try:
        future.result()
    except Exception as e:
        print(f"第{idx}行插入失败: {str(e)}")

print(f"插入完成,耗时: {(time.process_time() - start_insert)/60:.2f} 分钟")

2. 查询代码核心问题

  • 未获取全量数据:_current_rows仅返回当前分页的数据,Cassandra ResultSet是惰性加载的,需调用all()获取全部行。
  • COUNT查询用法错误:SELECT COUNT(*)返回单行结果,LIMIT 20000无效,且Cassandra的COUNT操作本身效率极低,小数据集直接遍历统计更可靠。

查询代码修复示例

query = f"SELECT * from {config['cassandra_db']['data_table']}"
simple_statement = SimpleStatement(query, consistency_level=ConsistencyLevel.ONE, fetch_size=1000)
execute_result = session.execute(simple_statement, timeout=None)
data = execute_result.all()  # 获取全量数据

print(f"查询到总行数: {len(data)}")
print(data.shape)

3. 额外排查步骤

  • 检查表的主键定义:确认define_columns中的主键是否合理,数据集是否存在重复主键值。
  • 查看DataStax控制台日志:确认是否有插入限流、主键冲突等错误记录。
  • 小批量测试:先插入100行验证全量插入/查询是否正常,逐步扩大范围定位问题。

内容的提问来源于stack exchange,提问作者Manu Vats

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.25 23:24:37