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

如何通过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.

需求1:通过Java循环生成并插入300,000条记录到Cassandra

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.UNLOGGED for 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.
需求2:读取文件中的300,000条Insert语句并插入Cassandra

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: BufferedReader reads 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_ONE for faster writes if eventual consistency is acceptable) using setConsistencyLevel(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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.22 08:10:06