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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.19 23:18:14