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

Java中是否有对应Cassandra cqlsh COPY FROM命令的实现方案?

Java批量加载CSV到Cassandra的高性能实现方案

你已经通过COPY Keyspace1.table1 FROM 'C:\cassandradata\priceinfo.csv' WITH DELIMITER = ',' and HEADER = true;实现了CSV批量导入,以下是几种用Java实现相近性能的方案:

一、DataStax官方Java驱动(性能最接近COPY FROM)

这是Cassandra官方推荐的Java客户端,底层实现和COPY FROM同源,性能最优。

步骤与代码示例

  1. 引入依赖(Maven)
<dependency>
    <groupId>com.datastax.oss</groupId>
    <artifactId>java-driver-core</artifactId>
    <version>4.17.0</version>
</dependency>
<dependency>
    <groupId>com.opencsv</groupId>
    <artifactId>opencsv</artifactId>
    <version>5.6</version>
</dependency>
  1. 实现CSV批量导入
import com.datastax.oss.driver.api.core.CqlSession;
import com.datastax.oss.driver.api.core.cql.AsyncResultSet;
import com.datastax.oss.driver.api.core.cql.PreparedStatement;
import com.opencsv.CSVReader;
import java.io.FileReader;
import java.util.concurrent.CompletableFuture;

public class CassandraCsvLoader {
    private static final String INSERT_QUERY = "INSERT INTO keyspace1.table1 (col1, col2, col3) VALUES (?, ?, ?)";
    private static final int BATCH_SIZE = 2000; // 根据数据大小调整,建议1000-5000

    public static void main(String[] args) {
        try (CqlSession session = CqlSession.builder().build();
             CSVReader reader = new CSVReader(new FileReader("C:\\cassandradata\\priceinfo.csv"))) {

            // 跳过表头
            reader.readNext();
            PreparedStatement stmt = session.prepare(INSERT_QUERY);
            int count = 0;
            CompletableFuture<AsyncResultSet> lastFuture = null;

            String[] nextLine;
            while ((nextLine = reader.readNext()) != null) {
                // 异步执行插入
                lastFuture = session.executeAsync(
                        stmt.bind(nextLine[0], nextLine[1], Double.parseDouble(nextLine[2])));
                
                count++;
                if (count % BATCH_SIZE == 0) {
                    // 每批量等待完成,避免内存溢出
                    lastFuture.join();
                    System.out.println("已插入 " + count + " 条数据");
                }
            }
            // 处理剩余数据
            if (lastFuture != null && count % BATCH_SIZE != 0) {
                lastFuture.join();
            }
            System.out.println("数据导入完成,共插入 " + count + " 条");
        } catch (Exception e) {
            e.printStackTrace();
        }
    }
}

性能优化点

  • 调整BATCH_SIZE:根据单条数据大小和Cassandra集群配置调整,避免单批次过大导致超时
  • 使用异步执行executeAsync:充分利用IO等待时间,提高吞吐量
  • 复用PreparedStatement:避免重复解析CQL语句,减少集群开销
  • 设置合适的一致性级别:如果业务允许,使用LOCAL_ONE替代默认的LOCAL_QUORUM,降低写入延迟

二、JDBC驱动方案

如果依赖JDBC生态,可以使用DataStax提供的JDBC驱动实现批量导入。

步骤与代码示例

  1. 引入依赖(Maven)
<dependency>
    <groupId>com.datastax.oss</groupId>
    <artifactId>jdbc-driver</artifactId>
    <version>1.4.1</version>
</dependency>
<dependency>
    <groupId>com.opencsv</groupId>
    <artifactId>opencsv</artifactId>
    <version>5.6</version>
</dependency>
  1. 实现批量导入
import java.sql.Connection;
import java.sql.DriverManager;
import java.sql.PreparedStatement;
import com.opencsv.CSVReader;
import java.io.FileReader;

public class CassandraJdbcCsvLoader {
    private static final String INSERT_QUERY = "INSERT INTO keyspace1.table1 (col1, col2, col3) VALUES (?, ?, ?)";
    private static final int BATCH_SIZE = 2000;

    public static void main(String[] args) {
        String url = "jdbc:cassandra://localhost:9042/keyspace1";
        try (Connection conn = DriverManager.getConnection(url, "", "");
             PreparedStatement stmt = conn.prepareStatement(INSERT_QUERY);
             CSVReader reader = new CSVReader(new FileReader("C:\\cassandradata\\priceinfo.csv"))) {

            reader.readNext(); // 跳过表头
            int count = 0;
            String[] nextLine;
            while ((nextLine = reader.readNext()) != null) {
                stmt.setString(1, nextLine[0]);
                stmt.setString(2, nextLine[1]);
                stmt.setDouble(3, Double.parseDouble(nextLine[2]));
                stmt.addBatch();
                
                count++;
                if (count % BATCH_SIZE == 0) {
                    stmt.executeBatch();
                    stmt.clearBatch();
                    System.out.println("已插入 " + count + " 条数据");
                }
            }
            // 处理剩余数据
            if (count % BATCH_SIZE != 0) {
                stmt.executeBatch();
            }
            System.out.println("数据导入完成,共插入 " + count + " 条");
        } catch (Exception e) {
            e.printStackTrace();
        }
    }
}

注意事项

  • JDBC的executeBatch性能略低于官方驱动的异步批量,适合已有JDBC代码的场景
  • 需注意Cassandra JDBC驱动的版本兼容性,建议和集群版本匹配

三、Spring Data Cassandra方案

如果使用Spring生态,Spring Data Cassandra提供了简化的批量操作API。

步骤与代码示例

  1. 引入依赖(Maven)
<dependency>
    <groupId>org.springframework.boot</groupId>
    <artifactId>spring-boot-starter-data-cassandra</artifactId>
</dependency>
<dependency>
    <groupId>com.opencsv</groupId>
    <artifactId>opencsv</artifactId>
    <version>5.6</version>
</dependency>
  1. 配置Cassandra连接(application.properties)
spring.data.cassandra.contact-points=localhost
spring.data.cassandra.port=9042
spring.data.cassandra.keyspace-name=keyspace1
  1. 定义实体类
import org.springframework.data.cassandra.core.mapping.PrimaryKey;
import org.springframework.data.cassandra.core.mapping.Table;

@Table("table1")
public class PriceInfo {
    @PrimaryKey
    private String col1;
    private String col2;
    private Double col3;

    // 构造器、getter、setter
    public PriceInfo(String col1, String col2, Double col3) {
        this.col1 = col1;
        this.col2 = col2;
        this.col3 = col3;
    }
    // 省略getter和setter
}
  1. 实现批量导入
import org.springframework.data.cassandra.core.CassandraTemplate;
import com.opencsv.CSVReader;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.CommandLineRunner;
import org.springframework.stereotype.Component;
import java.io.FileReader;
import java.util.ArrayList;
import java.util.List;

@Component
public class SpringDataCassandraLoader implements CommandLineRunner {
    @Autowired
    private CassandraTemplate cassandraTemplate;
    private static final int BATCH_SIZE = 2000;

    @Override
    public void run(String... args) throws Exception {
        try (CSVReader reader = new CSVReader(new FileReader("C:\\cassandradata\\priceinfo.csv"))) {
            reader.readNext(); // 跳过表头
            List<PriceInfo> batch = new ArrayList<>(BATCH_SIZE);
            int count = 0;
            String[] nextLine;

            while ((nextLine = reader.readNext()) != null) {
                PriceInfo info = new PriceInfo(nextLine[0], nextLine[1], Double.parseDouble(nextLine[2]));
                batch.add(info);
                count++;

                if (batch.size() == BATCH_SIZE) {
                    cassandraTemplate.insert(batch);
                    batch.clear();
                    System.out.println("已插入 " + count + " 条数据");
                }
            }
            // 处理剩余数据
            if (!batch.isEmpty()) {
                cassandraTemplate.insert(batch);
            }
            System.out.println("数据导入完成,共插入 " + count + " 条");
        }
    }
}

性能优化点

  • 开启异步插入:使用cassandraTemplate.insertAsync结合CompletableFuture批量处理
  • 调整Spring Data的批量配置:通过spring.data.cassandra.batch.size参数设置默认批量大小

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.25 16:36:11