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

多线程消息队列中Storage静态Map的keySet()返回空问题求助

多线程消息队列项目pubMsgs读取为空问题排查

问题现象

开发模拟消息队列的多线程项目时,定义了storage静态类(包含pubMsgs等静态Map字段),调用printPubMsgs方法时,pubMsgs.keySet()始终返回空集合,调试发现方法内pubMsgs大小为0,但确认此前已通过addMessage方法向其中添加元素,且其他线程请求storage时能正常返回数据。

核心代码片段

初始storage类片段

class storage{
    private static Map<String, List<String>> subLists = new HashMap<>();
    private static Map<String, List<Request>> pubMsgs = new HashMap<>();
    private static Lock lock = new ReentrantLock();

    public static void printPubMsgs(){
        lock.lock();
        Set<String> keyss = subLists.keySet();
        System.out.println("keyss:" + keyss);
        Set<String> keys = pubMsgs.keySet();
        System.out.println("keys:" + keys);
        for (String key :
                keys) {
            List<Request> list = pubMsgs.get(key);
            System.out.println("list1:" + list.toString());
            Iterator<Request> iterator = list.listIterator();
            while (iterator.hasNext()) {
                System.out.println(iterator);
                System.out.println("here");
            }
        }
        lock.unlock();
    }
}

addMessage方法

public static void addMessage(Request request) {
    lock.lock();
    if (pubMsgs.containsKey(request.getSource())) {
        pubMsgs.get(request.getSource()).add(request);
    } else {
        List<Request> list = new ArrayList<>();
        list.add(request);
        pubMsgs.put(request.getSource(), list);
    }
    lock.unlock();
}

项目核心类代码

Publisher.java

@Override
public void run() {
    try (Socket socket = new Socket(address, port);
        ObjectOutputStream oos = new ObjectOutputStream(socket.getOutputStream());) {

        oos.writeObject(new Request(pubHead, msg, id));
    } catch (Exception e) {
        e.printStackTrace();
    }
}

Subscriber.java

@Override
public void run() {
    try (Socket s = new Socket(address, port);
        ObjectOutputStream oos = new ObjectOutputStream(s.getOutputStream());) {
        oos.writeObject(new Request(subHead, sub_id, id));
        oos.close();
    } catch (Exception e) {
        e.printStackTrace();
    }

    try (Socket socket = new Socket(address, port);
        ObjectOutputStream oos1 = new ObjectOutputStream(socket.getOutputStream());) {
        oos1.writeObject(new Request(getHead, "", id));
//            sleep(1000);
        while (true) {
            ObjectInputStream ooi = new ObjectInputStream(socket.getInputStream());
            Request req = (Request) ooi.readObject();
            if (Objects.equals(req.getHead(), wrongHead)) {
                System.out.println(req.getBody());
                break;
            } else if (Objects.equals(req.getHead(), completeHead)) {
                System.out.println(req.getBody());
                break;
            }
            else{
                System.out.println(req.getBody());
            }
        }

    } catch (Exception e) {
        e.printStackTrace();
    }
}

Topic.java

@Override
public void run() {
    storage.update();
    try(ObjectInputStream ois = new ObjectInputStream(socket.getInputStream());){
        Request req = (Request) ois.readObject();
        if (req.getHead().equals(pubHead)) {
            System.out.println(">>> Successfully receive publish request");
            storage.addMessage(req);
        } else if (req.getHead().equals(getHead)) {
            System.out.println(">>> Successfully receive get request");
            List<Request> res = storage.getMessage(req.getSource());
            if (res == null) {
                ObjectOutputStream oos = new ObjectOutputStream(socket.getOutputStream());
                oos.writeObject(new Request(wrongHead, "you should subscribe first", "topic"));
            } else if (res.isEmpty()) {
                ObjectOutputStream oos = new ObjectOutputStream(socket.getOutputStream());
                oos.writeObject(new Request(wrongHead, "there are no messages from whom you subscribe", "topic"));
            } else {
                for (Request re : res) {
                    ObjectOutputStream oos = new ObjectOutputStream(socket.getOutputStream());
                    oos.writeObject(re);
                }
                ObjectOutputStream oos = new ObjectOutputStream(socket.getOutputStream());
                oos.writeObject(new Request(completeHead, ">>> all messages have been gotten", "topic"));
            }
        } else if (req.getHead().equals(subHead)) {
            storage.accessSub(req.getSource(), req.getBody());
            System.out.println(">>> Successfully receive subscribe request");
        }
    }catch (Exception e){
        e.printStackTrace();
    }
}

完整storage.java代码

class storage{
    private static Map<String, List<String>> subLists = new HashMap<>();
    private static Map<String, List<Request>> pubMsgs = new HashMap<>();
    private static Lock lock = new ReentrantLock();

    public static void printPubMsgs(){
        lock.lock();
        Set<String> keyss = subLists.keySet();
        System.out.println("keyss:" + keyss);
        Set<String> keys = pubMsgs.keySet();
        System.out.println("keys:" + keys);
        for (String key :
                keys) {
            List<Request> list = pubMsgs.get(key);
            System.out.println("list1:" + list.toString());
            Iterator<Request> iterator = list.listIterator();
            while (iterator.hasNext()) {
                System.out.println(iterator);
                System.out.println("here");
            }
        }
        lock.unlock();
    }

    public static void addMessage(Request request) {
        lock.lock();
        if (pubMsgs.containsKey(request.getSource())) {
            pubMsgs.get(request.getSource()).add(request);
        } else {
            List<Request> list = new ArrayList<>();
            list.add(request);
            pubMsgs.put(request.getSource(), list);
        }
        lock.unlock();
    }

    public static List<Request> getMessage(String source) {
        lock.lock();
        List<Request> result = new ArrayList<>();
        if (subLists.containsKey(source)) {
            List<String> list = subLists.get(source);
            for (String s : list) {
                if (!pubMsgs.containsKey(s)) continue;
                for (Request req : pubMsgs.get(s)) {
                    result.add(new Request(req.getHead(), req.getBody(), req.getSource()));
                }
            }
        } else {
            lock.unlock();
            return null; // wrong request
        }
        lock.unlock();
        return result; // if empty --> wrong request
    }

    public static void accessSub(String source, String destination) {
        lock.lock();
        if (subLists.containsKey(source)) {
            subLists.get(source).add(destination);
        } else {
            List<String> list = new ArrayList<>();
            list.add(destination);
            subLists.put(source, list);
        }
        lock.unlock();
    }

    /**
     * if the existence exceed 1000s, then remove
     */
    public static void update() {
        lock.lock();
        Set<String> keys = pubMsgs.keySet();
        for (String key : keys) {
            List<Request> list = pubMsgs.get(key);
            Iterator<Request> iterator = list.listIterator();
            while (iterator.hasNext()) {
                Date now = new Date();
                Date date = iterator.next().getDate();
                if ((now.getTime() - date.getTime()) / 1000 > 1000) {
                    iterator.remove();
                    System.out.println("test");
                }
            }
        }
        lock.unlock();
    }
}

可能的排查方向

1. 类加载器导致的静态实例不共享

如果调用printPubMsgs的线程与Topic线程的类加载器不同,会导致storage类被加载多次,各自拥有独立的静态字段实例。此时addMessage操作的是一个pubMsgs,而printPubMsgs读取的是另一个空的pubMsgs。

验证方式:在addMessage和printPubMsgs中添加日志,输出pubMsgs对象的哈希值:

// 在addMessage中
System.out.println("addMessage pubMsgs hash: " + pubMsgs.hashCode());
// 在printPubMsgs中
System.out.println("printPubMsgs pubMsgs hash: " + pubMsgs.hashCode());

如果两个哈希值不同,说明是类加载器问题。

2. Request.getSource()返回值不一致

检查Request类的getSource()方法,确认发布消息时传入的id(即Request的source)与printPubMsgs读取的keySet是否匹配。比如是否存在字符串大小写、空格或编码问题,导致addMessage存入的key与预期不符。

验证方式:在addMessage中添加日志,输出存入的key:

System.out.println("Adding message with source: " + request.getSource());

对比printPubMsgs中输出的keys,看是否存在对应的值。

3. 调用时机问题

确认printPubMsgs的调用是否在addMessage执行完成之后。如果printPubMsgs在消息添加前就被调用,自然会读取到空集合。

验证方式:在addMessage执行完成后添加标记,或者用同步机制确保printPubMsgs在消息添加后执行。

4. 内存可见性问题(虽用锁但仍需确认)

虽然ReentrantLock的lock/unlock操作会保证内存可见性,但如果存在代码路径未正确获取锁(比如pubMsgs被其他未加锁的代码修改),可能导致线程看不到最新值。检查所有操作pubMsgs的代码,确保都通过lock保护。

5. update方法的意外清理

Topic线程启动时会调用storage.update(),该方法会删除超过1000秒的消息。如果消息的getDate()返回的是未来时间,或者时间计算逻辑错误,可能导致刚添加的消息被误删。

验证方式:在update方法中添加日志,输出删除的消息信息:

if ((now.getTime() - date.getTime()) / 1000 > 1000) {
    System.out.println("Deleting message: " + iterator.next() + ", time diff: " + (now.getTime() - date.getTime())/1000);
    iterator.remove();
    System.out.println("test");
}

确认是否有刚添加的消息被误删。


内容的提问来源于stack exchange,提问作者EdgeRunner

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.19 17:11:59