如何使用Python将.csv格式文件上传导入至Cassandra数据库
Python实现CSV文件导入Cassandra数据库完整方案
前置准备
首先安装必要依赖,执行如下命令:pip install cassandra-driver
同时确认以下环境准备完成:
- Cassandra服务正常运行,网络可达
- 已经创建好对应keyspace,且表结构字段、类型和CSV文件字段完全匹配
- CSV文件编码统一,不存在乱码、缺列等格式问题
核心实现流程
- 初始化Cassandra连接,关联目标keyspace
- 预编译插入语句,避免重复解析SQL提升性能
- 逐行读取CSV文件,对字段做类型转换适配Cassandra表结构
- 按指定批量大小攒批插入,减少IO请求提升导入效率
- 所有数据导入完成后释放数据库连接
可直接落地的代码实现
请提前将代码中的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
相关产品推荐
相关产品推荐

