Python脚本执行Cassandra COPY命令时遇缓冲区超限断言错误
问题:Cassandra导入大TSV文件时出现序列化缓冲区溢出错误
尝试将500MB左右的制表符分隔(\t)数据文件通过Python调用CQL的COPY命令导入Cassandra时,触发以下错误:
Java.lang.RuntimeException: java.lang.AssertionError: Attempted serializing to buffer exceeded maximum of 65535 bytes: 1700053
相关Python代码:
def pd_to_cassandra_type(pd_series): if pd.api.types.is_integer_dtype(pd_series): return 'bigint' elif pd.api.types.is_float_dtype(pd_series): return 'double' else: return 'text' def clean_column_name(name): name = re.sub(r'\W+', '_', name) if name[0].isdigit(): name = '_' + name return name def create_table(session, keyspace_name, table_name, column_definition): session.execute(f""" CREATE KEYSPACE IF NOT EXISTS {keyspace_name} WITH REPLICATION = {{ 'class': 'SimpleStrategy', 'replication_factor': 1 }} """) session.set_keyspace(keyspace_name) session.execute(f""" CREATE TABLE IF NOT EXISTS {table_name} ( {column_definition}, id UUID PRIMARY KEY ) """) def split_and_load_data(): file_path = 'path/to/your/large_file.tab' keyspace_name = 'my_keyspace' table_name = "my_table" data = pd.read_csv(file_path, delimiter='\t') # Clean and reassign column names data.columns = [clean_column_name(col) for col in data.columns] user = "user_name" password = "1234" auth_provider = PlainTextAuthProvider(username=user, password=password) cluster = Cluster(["192.12.4.5"], auth_provider=auth_provider, connect_timeout=1000, control_connection_timeout=1000) session = cluster.connect() column_definitions = ", ".join([f"{col}" text for col in data.columns] # 根据第一列组创建表 create_table(session, keyspace_name, table_name, column_definitions) copy_command = f"cqlsh {"192.12.4.5"} -e \"COPY {keyspace}.{table} ({', '.join(data.columns)}) FROM '{file_path}' WITH DELIMITER='\t' AND HEADER=TRUE;\"" subprocess.run(copy_command, shell=True, check=True)
解决方案
1. 问题根源
这个错误是因为Cassandra的COPY命令默认限制单行数据序列化后的缓冲区大小为65535字节(64KB),而你的数据文件中存在单行数据(含所有字段)超过这个阈值的记录。
2. 调整COPY命令参数
直接修改COPY命令,增大允许的请求缓冲区大小,并调整批次大小:
# 修正后的copy_command copy_command = f"""cqlsh 192.12.4.5 -u {user} -p {password} -e "COPY {keyspace_name}.{table_name} ({', '.join(data.columns)}) FROM '{file_path}' WITH DELIMITER='\t' AND HEADER=TRUE AND MAX_REQUEST_SIZE=1048576 AND CHUNKSIZE=100;" """
MAX_REQUEST_SIZE=1048576:将单请求的最大缓冲区设为1MB(可根据实际单行大小调整,最大支持16MB)CHUNKSIZE=100:减小每次批量导入的行数,降低单批次内存压力- 补充了用户名密码参数,避免手动输入认证信息
3. 修正代码中的语法错误
原代码存在两处明显语法错误,会导致表创建失败:
column_definitions生成时缺少字段类型,且语法错误,应改为:
column_definitions = ", ".join([f"{col} {pd_to_cassandra_type(data[col])}" for col in data.columns])
- COPY命令中的变量名错误,
{keyspace}和{table}应替换为{keyspace_name}和{table_name}
4. 极端情况处理
如果调整参数后仍然报错,说明存在超大字段(比如单个text字段超过1MB):
- 若业务允许,可对超大文本字段进行截断处理
- 若需要保留完整数据,可将该字段拆分到单独的关联表中,通过主键关联主表
内容的提问来源于stack exchange,提问作者zehra özdemir
相关产品推荐
相关产品推荐

