求Python/Ruby将CSV导入Cassandra 3.11.3生产集群的可靠代码方案
推荐方案:Python + Cassandra Python Driver
针对你的生产场景,我强烈推荐使用Python来完成CSV到Cassandra的导入工作,相比Ruby,Python在这个场景下有明显优势:
- Python的
csv标准库对带特殊字符、换行的字段支持非常成熟(默认就能正确解析用引号包裹的多行字段) - 处理编码问题的工具链更完善(比如
chardet可以自动检测文件编码,避免UTF-8/GBK等编码兼容问题) - Cassandra官方维护的
python-driver稳定性强,生产环境案例丰富,对批量插入、大文本字段的处理更可靠 - 生态中有大量数据处理工具(如
pandas)可作为备选,应对复杂数据清洗需求
生产级Python示例代码
第一步:安装依赖
pip install cassandra-driver chardet
第二步:导入脚本(带完整错误处理、日志、批量插入)
import csv import chardet import logging from cassandra.cluster import Cluster from cassandra.auth import PlainTextAuthProvider from cassandra.query import BatchStatement, ConsistencyLevel from cassandra import OperationTimedOut, Unavailable # 配置生产环境日志 logging.basicConfig( level=logging.INFO, format='%(asctime)s - %(levelname)s - %(message)s', handlers=[ logging.FileHandler('csv_import.log'), logging.StreamHandler() ] ) logger = logging.getLogger(__name__) def detect_file_encoding(file_path): """自动检测CSV文件编码,解决UTF/GBK等编码不兼容问题""" with open(file_path, 'rb') as f: result = chardet.detect(f.read(10000)) # 读取前10KB检测编码 return result['encoding'] or 'utf-8' def connect_cassandra(node_ips, username=None, password=None): """建立Cassandra集群连接,生产环境建议配置认证""" auth_provider = PlainTextAuthProvider(username=username, password=password) if username else None try: cluster = Cluster( node_ips, consistency_level=ConsistencyLevel.LOCAL_QUORUM, # 7节点集群建议用LOCAL_QUORUM保证数据一致性 connect_timeout=30, idle_heartbeat_interval=30 ) session = cluster.connect('your_keyspace_name') # 替换为你的keyspace logger.info(f"成功连接Cassandra集群:{node_ips}") return cluster, session except Exception as e: logger.error(f"连接Cassandra失败:{str(e)}") raise def import_csv_to_cassandra(csv_path, session, table_name, batch_size=200): """批量导入CSV到Cassandra,处理特殊字符、换行、大文本""" encoding = detect_file_encoding(csv_path) logger.info(f"检测到CSV文件编码:{encoding}") # 准备插入语句(替换为你的表字段,用占位符避免SQL注入) insert_stmt = session.prepare(f""" INSERT INTO {table_name} (ticket_id, title, description, create_time, status) VALUES (?, ?, ?, ?, ?) """) insert_stmt.consistency_level = ConsistencyLevel.LOCAL_QUORUM processed_rows = 0 failed_rows = 0 with open(csv_path, 'r', encoding=encoding, newline='') as csvfile: # 使用csv.DictReader自动处理带引号的换行字段 reader = csv.DictReader(csvfile) batch = BatchStatement(consistency_level=ConsistencyLevel.LOCAL_QUORUM) for row_num, row in enumerate(reader, start=2): # 从第2行开始,跳过表头 try: # 处理空值(Cassandra不允许插入None,根据表结构替换为默认值或跳过) ticket_id = row.get('ticket_id') or '' title = row.get('title') or '' description = row.get('description') or '' # 大文本直接传入,Cassandra text类型支持 create_time = row.get('create_time') or '' status = row.get('status') or '' batch.add(insert_stmt, (ticket_id, title, description, create_time, status)) # 达到批量大小则执行插入 if len(batch) >= batch_size: session.execute(batch) processed_rows += len(batch) logger.info(f"已处理 {processed_rows} 条记录") batch.clear() except (OperationTimedOut, Unavailable) as e: # 处理Cassandra连接超时或不可用,重试当前批次 logger.error(f"批次插入超时/集群不可用,行号:{row_num},错误:{str(e)}") batch.clear() failed_rows += 1 except Exception as e: # 处理其他错误(如数据类型不匹配),记录并跳过该行 logger.error(f"处理行失败,行号:{row_num},内容:{row},错误:{str(e)}") failed_rows += 1 # 处理剩余的不足批量大小的记录 if len(batch) > 0: try: session.execute(batch) processed_rows += len(batch) logger.info(f"已处理 {processed_rows} 条记录(剩余批次)") except Exception as e: logger.error(f"剩余批次插入失败:{str(e)}") failed_rows += len(batch) logger.info(f"导入完成!成功处理 {processed_rows} 条,失败 {failed_rows} 条") if __name__ == "__main__": # 配置参数(替换为你的生产环境信息) CASSANDRA_NODES = ['node1-ip', 'node2-ip', 'node3-ip', 'node4-ip', 'node5-ip', 'node6-ip', 'node7-ip'] CASSANDRA_USER = 'your_username' # 如果集群开启认证 CASSANDRA_PASS = 'your_password' CSV_FILE_PATH = '/path/to/your/tickets.csv' CASSANDRA_TABLE = 'your_ticket_table' cluster, session = None, None try: cluster, session = connect_cassandra(CASSANDRA_NODES, CASSANDRA_USER, CASSANDRA_PASS) import_csv_to_cassandra(CSV_FILE_PATH, session, CASSANDRA_TABLE) except Exception as e: logger.error(f"导入任务失败:{str(e)}") finally: # 确保关闭连接 if session: session.shutdown() if cluster: cluster.shutdown() logger.info("Cassandra连接已关闭")
生产环境注意事项
- 批量大小调整:根据单条记录的大小调整
batch_size(建议100-500之间),避免Cassandra的batch过大导致性能问题 - 一致性级别:7节点集群使用
LOCAL_QUORUM可以保证数据在大多数节点写入成功,兼顾一致性和性能 - 编码处理:如果已知CSV编码,可直接指定编码跳过检测,提高效率
- 错误重试:对于超时等可恢复错误,可添加重试逻辑(比如用
tenacity库) - 监控:可以添加Prometheus指标或写入监控系统,跟踪导入进度和错误率
内容的提问来源于stack exchange,提问作者Hareesha
相关产品推荐
相关产品推荐

