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同源,性能最优。
步骤与代码示例
- 引入依赖(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>
- 实现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驱动实现批量导入。
步骤与代码示例
- 引入依赖(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>
- 实现批量导入
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。
步骤与代码示例
- 引入依赖(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>
- 配置Cassandra连接(application.properties)
spring.data.cassandra.contact-points=localhost spring.data.cassandra.port=9042 spring.data.cassandra.keyspace-name=keyspace1
- 定义实体类
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 }
- 实现批量导入
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
相关产品推荐
相关产品推荐

