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
相关产品推荐
相关产品推荐

