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

使用Java驱动异步插入Cassandra时高频写入丢失数据问题排查

解决Cassandra异步插入时的数据丢失问题

你的问题很典型——直接在循环里调用executeAsync而不做任何控制,确实会导致数据丢失,主要有几个核心原因和对应的解决办法,咱们一步步来梳理:

核心问题分析

  • 未等待异步操作完成:当main方法执行完循环后,程序可能直接退出,此时还有大量executeAsync的请求没完成,这些未完成的插入就会丢失。
  • 无限制提交请求压垮集群:500万条请求一下子全发出去,Cassandra的队列会溢出,导致部分请求被丢弃,而且你没有处理异步操作的异常,失败的请求也没重试。
  • 缺少请求追踪与错误处理:异步操作的失败不会直接抛出异常,如果你不主动处理ResultSetFuture的结果,根本不知道哪些请求失败了。

具体解决方案

1. 等待所有异步请求完成

不要让程序在循环结束后直接退出,要收集所有的ResultSetFuture,然后等待它们全部完成。示例代码:

public static void main(String[] args) {
    FileReader fr = null;
    List<ResultSetFuture> futures = new ArrayList<>();
    try {
        fr = new FileReader("the-file-name.txt");
        BufferedReader br = new BufferedReader(fr);
        String sCurrentLine;
        long time1 = System.currentTimeMillis();
        
        while ((sCurrentLine = br.readLine()) != null) {
            ResultSetFuture future = session.executeAsync(sCurrentLine);
            futures.add(future);
        }
        
        // 等待所有异步请求完成,同时处理失败情况
        for (ResultSetFuture future : futures) {
            try {
                future.get(); // 阻塞直到请求完成
            } catch (InterruptedException | ExecutionException e) {
                System.err.println("插入失败,语句可能丢失: " + e.getMessage());
                // 这里可以加入重试逻辑或者记录异常语句
            }
        }
        
        long time2 = System.currentTimeMillis();
        System.out.println("总耗时: " + (time2 - time1) + "ms");
    } catch (IOException e) {
        e.printStackTrace();
    } finally {
        if (fr != null) {
            try {
                fr.close();
            } catch (IOException e) {
                e.printStackTrace();
            }
        }
        if (session != null) {
            session.close();
        }
    }
}

2. 控制并发请求数

直接提交500万请求会压垮Cassandra,最好用Semaphore或者线程池来控制并发数,避免集群过载:

public static void main(String[] args) {
    FileReader fr = null;
    // 根据集群配置调整并发数,比如100个同时请求
    Semaphore semaphore = new Semaphore(100);
    ExecutorService executor = Executors.newFixedThreadPool(100);
    
    try {
        fr = new FileReader("the-file-name.txt");
        BufferedReader br = new BufferedReader(fr);
        String sCurrentLine;
        long time1 = System.currentTimeMillis();
        
        while ((sCurrentLine = br.readLine()) != null) {
            semaphore.acquire(); // 获取许可,不够则等待
            executor.submit(() -> {
                try {
                    ResultSetFuture future = session.executeAsync(sCurrentLine);
                    future.get();
                } catch (Exception e) {
                    System.err.println("插入失败: " + e.getMessage());
                    // 可加入重试逻辑,注意幂等性避免重复插入
                } finally {
                    semaphore.release(); // 释放许可,让下一个请求执行
                }
            });
        }
        
        executor.shutdown();
        executor.awaitTermination(1, TimeUnit.HOURS); // 等待所有任务完成
        
        long time2 = System.currentTimeMillis();
        System.out.println("总耗时: " + (time2 - time1) + "ms");
    } catch (IOException | InterruptedException e) {
        e.printStackTrace();
    } finally {
        if (fr != null) {
            try {
                fr.close();
            } catch (IOException e) {
                e.printStackTrace();
            }
        }
        if (session != null) {
            session.close();
        }
    }
}

3. 使用批量插入优化性能(可选)

如果你的INSERT语句针对同一张表,可以把多条语句合并成BATCH语句,减少请求次数,同时降低集群压力:

public static void main(String[] args) {
    FileReader fr = null;
    List<String> batchStatements = new ArrayList<>();
    int batchSize = 100; // 每100条合并为一个batch
    
    try {
        fr = new FileReader("the-file-name.txt");
        BufferedReader br = new BufferedReader(fr);
        String sCurrentLine;
        long time1 = System.currentTimeMillis();
        
        while ((sCurrentLine = br.readLine()) != null) {
            batchStatements.add(sCurrentLine);
            if (batchStatements.size() >= batchSize) {
                // 构建batch语句
                StringBuilder batch = new StringBuilder("BEGIN BATCH ");
                for (String stmt : batchStatements) {
                    batch.append(stmt).append("; ");
                }
                batch.append("APPLY BATCH;");
                
                session.executeAsync(batch.toString()).get(); // 等待batch完成
                batchStatements.clear();
            }
        }
        
        // 处理剩余的未达批量的语句
        if (!batchStatements.isEmpty()) {
            StringBuilder batch = new StringBuilder("BEGIN BATCH ");
            for (String stmt : batchStatements) {
                batch.append(stmt).append("; ");
            }
            batch.append("APPLY BATCH;");
            session.executeAsync(batch.toString()).get();
        }
        
        long time2 = System.currentTimeMillis();
        System.out.println("总耗时: " + (time2 - time1) + "ms");
    } catch (Exception e) {
        e.printStackTrace();
    } finally {
        if (fr != null) {
            try {
                fr.close();
            } catch (IOException e) {
                e.printStackTrace();
            }
        }
        if (session != null) {
            session.close();
        }
    }
}

4. 关键注意事项

  • 重试机制:对于失败的请求,一定要加入重试逻辑(注意幂等性,比如主键重复插入不会改变数据的情况,重试安全;否则要避免重复插入)。
  • 集群配置:确保Cassandra集群有足够资源,可调整concurrent_writes、write_request_timeout_in_ms等参数适配高并发写入。
  • 日志监控:开启驱动的DEBUG日志(比如com.datastax.driver.core包),追踪每个请求的状态,方便定位问题。

内容的提问来源于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:29:08