Java多线程KV存储高并发GET返回旧版本问题排查
并发请求下KV存储版本不一致问题
自研了一款带版本控制的Web服务器,KV存储基于嵌套HashMap实现。单请求场景正常,但发送25000次并发请求时,GET接口总是返回对应值的旧版本。已尝试用Lock和ConcurrentHashMap实现同步访问,问题仍未解决。
核心代码
DataManager类
package cis5550.kvs; import java.util.Map; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.locks.Lock; import java.util.concurrent.locks.ReentrantReadWriteLock; public class DataManager { private Map<String, Map<String, Map<String, Map<Integer, byte[]>>>> data; private ReentrantReadWriteLock lock; public DataManager() { data = new ConcurrentHashMap<>(); lock = new ReentrantReadWriteLock(); } public synchronized String put(String table, String row, String column, byte[] value) { try { lock.writeLock().lock(); Map<String, Map<String, Map<Integer, byte[]>>> rowMap = data.get(table); if (rowMap == null) { rowMap = new ConcurrentHashMap<>(); data.put(table, rowMap); } Map<String, Map<Integer, byte[]>> colMap = rowMap.get(row); if (colMap == null) { colMap = new ConcurrentHashMap<>(); rowMap.put(row, colMap); } Map<Integer, byte[]> versionMap = colMap.get(column); if (versionMap == null) { versionMap = new ConcurrentHashMap<>(); colMap.put(column, versionMap); } int latestVersion = getLatestVersion(versionMap); int newVersion = latestVersion + 1; versionMap.put(newVersion, value); return String.valueOf(newVersion); }finally { lock.writeLock().unlock(); } } private synchronized int getLatestVersion(Map<Integer, byte[]> versionMap) { return versionMap.keySet().stream().max(Integer::compareTo).orElse(0); } public synchronized byte[] get(String table, String row, String column, int version) { try { lock.readLock().lock(); Map<String, Map<String, Map<Integer, byte[]>>> rowMap = data.get(table); if (rowMap == null) { return null; } Map<String, Map<Integer, byte[]>> colMap = rowMap.get(row); if (colMap == null) { return null; } Map<Integer, byte[]> versionMap = colMap.get(column); if (versionMap == null) { return null; } return versionMap.get(version); }finally { lock.readLock().unlock(); } } public synchronized int getLatestVersion(String table, String row, String column) { Map<String, Map<String, Map<Integer, byte[]>>> rowMap = data.get(table); if (rowMap == null) { return 0; } Map<String, Map<Integer, byte[]>> colMap = rowMap.get(row); if (colMap == null) { return 0; } Map<Integer, byte[]> versionMap = colMap.get(column); if (versionMap == null || versionMap.isEmpty()) { return 0; } return getLatestVersion(versionMap); } }
Worker类
package cis5550.kvs; import cis5550.webserver.Server; import java.nio.charset.StandardCharsets; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; public class Worker extends cis5550.generic.Worker { private static final int MAX_THREADS = 1000; public static void main(String[] args) { if (args.length < 3) { System.out.println("Enter the required <port> <storage directory> <ip:port>"); System.exit(1); } Server.port(Integer.parseInt(args[0])); startPingThread(args[2], args[0], args[1]); DataManager dataManager = new DataManager(); ExecutorService threadPool = Executors.newFixedThreadPool(MAX_THREADS); Server.put("/data/:T/:R/:C", (req, res) -> { try { String tableName = req.params("T"); String rowName = req.params("R"); String columnName = req.params("C"); if (req.queryParams().contains("ifcolumn") && req.queryParams().contains("equals")) { String ifColumnName = req.queryParams("ifcolumn"); String ifColumnValue = req.queryParams("equals"); int latestVersion = dataManager.getLatestVersion(tableName, rowName, columnName); byte[] byteData = dataManager.get(tableName, rowName, ifColumnName , latestVersion) != null ? dataManager.get(tableName, rowName, ifColumnName , latestVersion) : new byte[0]; String data = new String(byteData, StandardCharsets.UTF_8); if (!data.equals("") && data.equals(ifColumnValue)) { threadPool.execute(() -> { res.header("version", dataManager.put(tableName, rowName, columnName, req.bodyAsBytes())); }); return "OK"; } else { return "FAIL"; } } else { threadPool.execute(() -> { res.header("version", dataManager.put(tableName, rowName, columnName, req.bodyAsBytes())); }); return "OK"; } } catch (Exception e) { res.status(404, "FAIL"); return null; } }); Server.get("/data/:T/:R/:C", (req, res) -> { try { String tableName = req.params("T"); String rowName = req.params("R"); String columnName = req.params("C"); if (req.queryParams().contains("version")) { int version = Integer.parseInt(req.queryParams("version")); String data = new String(dataManager.get(tableName, rowName, columnName, version), StandardCharsets.UTF_8); res.header("version", req.params("version")); res.body(data); } else { int latestVersion = dataManager.getLatestVersion(tableName, rowName, columnName); String data = new String(dataManager.get(tableName, rowName, columnName, latestVersion), StandardCharsets.UTF_8); res.header("version", String.valueOf(latestVersion)); res.body(data); } } catch (Exception e) { res.status(404, "FAIL"); } return null; }); } }
问题根源与修复方案
核心问题
- PUT请求异步执行导致数据不一致:Worker中PUT操作被丢给线程池异步执行,但直接返回了"OK"。客户端收到响应时,数据可能还未完成写入,后续GET请求会读到旧值。
- 同步机制冗余且存在漏洞:DataManager的方法同时使用
synchronized和ReentrantReadWriteLock,但getLatestVersion(String table, String row, String column)未加读锁,直接操作ConcurrentHashMap,并发场景下会读到不一致的版本信息。 - 版本号计算线程不安全:通过遍历keySet找最大版本号,ConcurrentHashMap的keySet是弱一致性的,并发修改时可能遍历到不完整的键集合,导致多个线程写入相同版本号。
修复步骤
1. 修正PUT请求的异步逻辑
PUT操作需等待写入完成后再响应,避免客户端提前收到结果:
// 替换原异步执行代码,改为同步执行 String version = dataManager.put(tableName, rowName, columnName, req.bodyAsBytes()); res.header("version", version);
若需保留异步,需确保响应在写入完成后发送(需确认Web服务器支持异步响应完成机制):
CompletableFuture.runAsync(() -> { String version = dataManager.put(tableName, rowName, columnName, req.bodyAsBytes()); res.header("version", version); }, threadPool);
2. 清理冗余同步,完善锁覆盖
移除方法上的synchronized,统一用读写锁覆盖所有操作:
public int getLatestVersion(String table, String row, String column) { try { lock.readLock().lock(); // 新增读锁 Map<String, Map<String, Map<Integer, byte[]>>> rowMap = data.get(table); if (rowMap == null) { return 0; } Map<String, Map<Integer, byte[]>> colMap = rowMap.get(row); if (colMap == null) { return 0; } Map<Integer, byte[]> versionMap = colMap.get(column); if (versionMap == null || versionMap.isEmpty()) { return 0; } return getLatestVersion(versionMap); } finally { lock.readLock().unlock(); // 释放读锁 } } // 移除其他方法上的synchronized修饰符,仅保留读写锁逻辑
3. 优化版本号存储,避免遍历计算
用AtomicInteger存储最新版本号,替代遍历keySet的方式:
// 新增内部类存储版本数据 private static class VersionedData { private final ConcurrentHashMap<Integer, byte[]> versions = new ConcurrentHashMap<>(); private final AtomicInteger latestVersion = new AtomicInteger(0); } // 修改DataManager的data结构 private Map<String, Map<String, Map<String, VersionedData>>> data; // 修改put方法 public String put(String table, String row, String column, byte[] value) { try { lock.writeLock().lock(); Map<String, Map<String, VersionedData>> rowMap = data.get(table); if (rowMap == null) { rowMap = new ConcurrentHashMap<>(); data.put(table, rowMap); } Map<String, VersionedData> colMap = rowMap.get(row); if (colMap == null) { colMap = new ConcurrentHashMap<>(); rowMap.put(row, colMap); } VersionedData versionedData = colMap.get(column); if (versionedData == null) { versionedData = new VersionedData(); colMap.put(column, versionedData); } int newVersion = versionedData.latestVersion.incrementAndGet(); versionedData.versions.put(newVersion, value); return String.valueOf(newVersion); } finally { lock.writeLock().unlock(); } }
内容的提问来源于stack exchange,提问作者Vamsi Konakanchi
相关产品推荐
相关产品推荐

