使用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
相关产品推荐
相关产品推荐

