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

求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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.13 07:32:28