如何通过Java向Cassandra插入30万条记录:循环与读文件场景
Hey there! Let's break down how to handle both of your Cassandra bulk insertion tasks with Java, including optimized code that's reliable for 300k records. I'll build out the partial code you provided into complete, working examples.
First, let's assume you have a target table in Cassandra (if not, create it first with a CQL statement). We'll use batch statements to group inserts (instead of single-row inserts, which are way too slow for 300k records) and leverage try-with-resources to safely manage resources like the Cluster and Session.
完整代码示例
import com.datastax.driver.core.*; import java.util.Random; public class CassandraBulkInsertFromLoop { // Configure your Cassandra connection details private static final String CONTACT_POINT = "127.0.0.1"; private static final String KEYSPACE = "test"; private static final String TABLE = "user_data"; // Replace with your actual table name private static final int TOTAL_RECORDS = 300000; private static final int BATCH_SIZE = 1000; // Adjust based on your cluster's capacity public static void main(String[] args) { // Use try-with-resources to auto-close Cluster and Session (avoids resource leaks) try (Cluster cluster = Cluster.builder().addContactPoint(CONTACT_POINT).build(); Session session = cluster.connect(KEYSPACE)) { // Prepare reusable INSERT statement (cuts down on Cassandra parsing overhead) String insertCql = String.format("INSERT INTO %s (id, name, email) VALUES (?, ?, ?)", TABLE); PreparedStatement preparedStmt = session.prepare(insertCql); Random random = new Random(); BatchStatement batch = new BatchStatement(BatchType.UNLOGGED); // Faster for independent records int recordCount = 0; for (int i = 1; i <= TOTAL_RECORDS; i++) { // Generate sample data (replace with your actual data logic) String name = "User_" + i; String email = "user_" + i + "@example.com"; // Add record to batch batch.add(preparedStmt.bind(i, name, email)); recordCount++; // Execute batch when we hit the batch size, or on the final iteration if (recordCount % BATCH_SIZE == 0 || i == TOTAL_RECORDS) { session.execute(batch); batch.clear(); // Reset for next group System.out.println("Inserted " + recordCount + " records so far..."); } } System.out.println("Successfully inserted all " + TOTAL_RECORDS + " records!"); } catch (Exception e) { System.err.println("Error during bulk insertion: " + e.getMessage()); e.printStackTrace(); } } }
关键优化说明
- Unlogged Batch: Use
BatchType.UNLOGGEDfor faster writes (logged batches add atomicity overhead that's unnecessary for independent records). - Prepared Statement: Reusing a prepared statement reduces repeated CQL parsing on Cassandra nodes.
- Batch Size: Adjust
BATCH_SIZE(500-2000 is a safe range) based on your cluster's memory and throughput—too large can cause timeouts. - Try-with-resources: Automatically cleans up connections to avoid memory leaks.
For this task, we'll read the file line by line, group valid INSERT statements into batches, and execute them. We'll add safeguards to skip malformed lines or comments to avoid breaking the entire job.
完整代码示例
import com.datastax.driver.core.*; import java.io.BufferedReader; import java.io.FileReader; import java.io.IOException; public class CassandraBulkInsertFromFile { private static final String CONTACT_POINT = "127.0.0.1"; private static final String KEYSPACE = "test"; private static final String FILE_PATH = "the-file-name.txt"; // Your file path private static final int BATCH_SIZE = 1000; public static void main(String[] args) { try (Cluster cluster = Cluster.builder().addContactPoint(CONTACT_POINT).build(); Session session = cluster.connect(KEYSPACE); BufferedReader br = new BufferedReader(new FileReader(FILE_PATH))) { BatchStatement batch = new BatchStatement(BatchType.UNLOGGED); int recordCount = 0; String line; while ((line = br.readLine()) != null) { line = line.trim(); // Skip empty lines, comments, or non-INSERT statements if (line.isEmpty() || line.startsWith("--") || !line.toUpperCase().startsWith("INSERT")) { if (!line.isEmpty() && !line.startsWith("--")) { System.err.println("Skipping invalid line: " + line); } continue; } // Add valid INSERT statement to batch batch.add(new SimpleStatement(line)); recordCount++; // Execute batch when batch size is reached if (recordCount % BATCH_SIZE == 0) { session.execute(batch); batch.clear(); System.out.println("Processed " + recordCount + " statements so far..."); } } // Execute any remaining statements in the final batch if (!batch.isEmpty()) { session.execute(batch); System.out.println("Processed final batch of " + batch.size() + " statements"); } System.out.println("Successfully processed all " + recordCount + " INSERT statements!"); } catch (IOException e) { System.err.println("Error reading file: " + e.getMessage()); e.printStackTrace(); } catch (Exception e) { System.err.println("Error executing CQL statements: " + e.getMessage()); e.printStackTrace(); } } }
注意事项
- Efficient File Reading:
BufferedReaderreads lines incrementally, so we don't load the entire 300k-line file into memory. - Validation Checks: Prevents invalid lines from crashing the entire insertion job.
- Batch Logic: Same batch size principles apply here—tune based on your cluster's performance.
通用最佳实践
- Cassandra Tuning: Ensure your cluster is configured for high write throughput (adjust
concurrent_writes,commitlog_segment_size_in_mb, and heap memory as needed). - Consistency Level: Set an appropriate consistency level (e.g.,
LOCAL_ONEfor faster writes if eventual consistency is acceptable) usingsetConsistencyLevel(ConsistencyLevel.LOCAL_ONE)on your statements. - Error Recovery: For production use, add retry logic for transient failures and log failed statements for later reprocessing.
内容的提问来源于stack exchange,提问作者Vickie Jack

