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

如何使用Python将.csv格式文件上传导入至Cassandra数据库

Python实现CSV文件导入Cassandra数据库完整方案

前置准备

首先安装必要依赖,执行如下命令:
pip install cassandra-driver
同时确认以下环境准备完成:

  • Cassandra服务正常运行,网络可达
  • 已经创建好对应keyspace,且表结构字段、类型和CSV文件字段完全匹配
  • CSV文件编码统一,不存在乱码、缺列等格式问题

核心实现流程

  1. 初始化Cassandra连接,关联目标keyspace
  2. 预编译插入语句,避免重复解析SQL提升性能
  3. 逐行读取CSV文件,对字段做类型转换适配Cassandra表结构
  4. 按指定批量大小攒批插入,减少IO请求提升导入效率
  5. 所有数据导入完成后释放数据库连接

可直接落地的代码实现

请提前将代码中的Cassandra连接配置、keyspace名称、表名、CSV路径、表字段替换为自身实际业务配置

from cassandra.cluster import Cluster
from cassandra.query import BatchStatement
import csv

# -------------------------- 配置项 start --------------------------
CASSANDRA_HOST = '127.0.0.1'  # Cassandra节点地址,集群可填多个地址如['192.168.1.1','192.168.1.2']
CASSANDRA_PORT = 9042  # Cassandra默认服务端口
TARGET_KEYSPACE = 'your_keyspace'  # 目标keyspace名称
TARGET_TABLE = 'your_table_name'  # 目标表名
CSV_FILE_PATH = './test_data.csv'  # 本地CSV文件路径
BATCH_INSERT_SIZE = 100  # 每多少条数据执行一次批量插入,建议50-200之间
# -------------------------- 配置项 end ----------------------------

# 初始化Cassandra连接
cluster = Cluster([CASSANDRA_HOST], port=CASSANDRA_PORT)
session = cluster.connect(TARGET_KEYSPACE)

# 预编译插入SQL,此处VALUES前的字段需要和你的表字段、CSV字段顺序完全一致
# 示例假设表字段为:id(int)、username(text)、age(int)、register_time(timestamp)
insert_stmt = session.prepare(f"""
    INSERT INTO {TARGET_TABLE} (id, username, age, register_time)
    VALUES (?, ?, ?, ?)
""")

if __name__ == '__main__':
    batch = BatchStatement()
    total_insert = 0

    with open(CSV_FILE_PATH, 'r', encoding='utf-8') as f:
        csv_reader = csv.reader(f)
        # 如果CSV第一行是表头,执行下一行跳过表头;没有表头则删除下面这行
        next(csv_reader)

        for row in csv_reader:
            # 此处根据你的字段类型做转换,示例为将CSV的字符串转成对应类型
            processed_data = [
                int(row[0]),  # id字段转int
                row[1],  # username字段保持字符串
                int(row[2]),  # age字段转int
                row[3]  # 时间字段如果是标准格式字符串Cassandra可自动识别,也可以手动转datetime类型
            ]
            batch.add(insert_stmt, processed_data)
            total_insert += 1

            # 达到批量阈值执行插入
            if total_insert % BATCH_INSERT_SIZE == 0:
                session.execute(batch)
                batch.clear()
                print(f"已完成{total_insert}条数据导入")

        # 插入剩余不足批量阈值的数据
        if total_insert % BATCH_INSERT_SIZE != 0:
            session.execute(batch)
            print(f"全部数据导入完成,累计导入{total_insert}条")

    # 释放连接
    session.shutdown()
    cluster.shutdown()

注意事项

  • 数据类型匹配:CSV读取的所有字段默认都是字符串类型,必须转换为和Cassandra表字段对应的类型,否则会抛出插入异常
  • 批量大小控制:批量插入阈值不要超过200,过大会导致Cassandra节点写入压力过高触发超时
  • 主键重复处理:Cassandra遇到主键重复的数据会默认覆盖旧数据,有去重需求的场景需要提前对CSV数据做主键去重
  • 大文件优化:如果CSV文件超过1G,建议增加断点续传逻辑,记录已导入的行数,异常中断后不需要从头开始导入

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.06 01:24:03