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

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;
        });
    }
}

问题根源与修复方案

核心问题

  1. PUT请求异步执行导致数据不一致:Worker中PUT操作被丢给线程池异步执行,但直接返回了"OK"。客户端收到响应时,数据可能还未完成写入,后续GET请求会读到旧值。
  2. 同步机制冗余且存在漏洞:DataManager的方法同时使用synchronized和ReentrantReadWriteLock,但getLatestVersion(String table, String row, String column)未加读锁,直接操作ConcurrentHashMap,并发场景下会读到不一致的版本信息。
  3. 版本号计算线程不安全:通过遍历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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.30 21:18:49