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

ConcurrentHashMap+MutableInteger并发计数结果异常问题排查

并发计数错误问题排查与修复

问题场景

使用ConcurrentHashMap存储时间维度(日/时/分)的流量计数,键为时间字符串,值为自定义的MutableInteger(其getValue和setValue方法已加synchronized锁),但最终计数结果远低于预期的1000次,计数丢失严重。

相关代码如下:

MutableInteger类

public class MutableInteger {
    private volatile int value;

    public MutableInteger() {}

    public synchronized int getValue() {
        return value;
    }

    public synchronized void setValue(int value) {
        this.value = value;
    }
}

主逻辑类

import java.text.SimpleDateFormat;
import java.util.Calendar;
import java.util.concurrent.*;

public class LinkedQueue {

    static ConcurrentHashMap<String,MutableInteger> map = new ConcurrentHashMap();

    public static class CountQueue{
        static BlockingQueue<Calendar> queue = new LinkedBlockingQueue<Calendar>();

        public void produce() throws InterruptedException{
            queue.put(Calendar.getInstance());
        }
        public Calendar consume() throws InterruptedException{
            return queue.take();
        }
    }

    public static void recordFlow(){
        final CountQueue queue = new CountQueue();
        class Producer implements Runnable {
            public void run() {
                try {
                    queue.produce();
                } catch (Exception ex) {
                    ex.printStackTrace();
                }
            }
        }
        class Consumer implements Runnable {
            public void run() {
                try {
                    getdata(queue.consume());
                } catch (Exception ex) {
                    ex.printStackTrace();
                }
            }
        }

        ExecutorService service = Executors.newCachedThreadPool();
        Producer producer = new Producer();
        Consumer consumer = new Consumer();
        service.submit(producer);
        service.submit(consumer);
        service.shutdown();
    }

    public static void getdata(Calendar c){
        try{
            saveFlow(c);
        }catch (Exception e){
            e.printStackTrace();
        }
    }

    // 接口流量记录
    public static void saveFlow(Calendar c){
        System.out.println("Thread Name:"+Thread.currentThread().getName());
        SimpleDateFormat sdfDate = new SimpleDateFormat("yyyy-MM-dd");
        SimpleDateFormat sdfHour = new SimpleDateFormat("yyyy-MM-dd HH");
        SimpleDateFormat sdfMinute = new SimpleDateFormat("yyyy-MM-dd HH:mm");

        String dayAcceptKey = sdfDate.format(c.getTime());
        String hourAcceptKey = sdfHour.format(c.getTime());
        String minuteAcceptKey = sdfMinute.format(c.getTime());
        
        MutableInteger newMinuteAcceptKey = new MutableInteger();
        newMinuteAcceptKey.setValue(1);
        MutableInteger oldMinuteAcceptKey = map.put(minuteAcceptKey, newMinuteAcceptKey);
        if(null != oldMinuteAcceptKey) {
            newMinuteAcceptKey.setValue(oldMinuteAcceptKey.getValue() + 1);
        }
        System.out.println("minute in map:"+ minuteAcceptKey + ","+ newMinuteAcceptKey.getValue());
        MutableInteger newHourAcceptKey = new MutableInteger();
        newHourAcceptKey.setValue(1);
        MutableInteger oldHourAcceptKey = map.put(hourAcceptKey, newHourAcceptKey);
        if(null != oldHourAcceptKey) {
            newHourAcceptKey.setValue(oldHourAcceptKey.getValue() + 1);
        }
        System.out.println("hour in map:" + hourAcceptKey + "," + newHourAcceptKey.getValue());
        MutableInteger newDayAcceptKey = new MutableInteger();
        newDayAcceptKey.setValue(1);
        MutableInteger oldDayAcceptKey = map.put(dayAcceptKey, newDayAcceptKey);
        if(null != oldDayAcceptKey) {
            newDayAcceptKey.setValue(oldDayAcceptKey.getValue() + 1);
        }
        System.out.println("day in map:" + dayAcceptKey + "," + newDayAcceptKey.getValue());
    }

    public static void main(String[] args){
        for(int i=0;i<1000;i++) {
            recordFlow();
        }
        try {
            Thread.sleep(2000);
        } catch (InterruptedException e) {
            e.printStackTrace();
        }
        System.out.println("==========================");
        for(String key:map.keySet()){
            System.out.println("key and value:" +key + " " + map.get(key).getValue());
        }
    }
}

问题原因分析

  1. 核心逻辑非原子性:
    ConcurrentHashMap的put方法是原子操作,但整个计数流程是"创建新对象→put到map→获取旧对象→更新新对象值",这一系列操作不是原子的。多个线程同时处理同一个时间key时,会发生覆盖:

    • 线程A创建值为1的MutableInteger,并put到map;
    • 线程B同时创建值为1的MutableInteger,覆盖线程A的对象;
    • 线程A拿到旧值(可能是更早的对象),更新自己的新对象,但这个对象已经不在map里了,map里是线程B的对象,导致线程A的计数丢失。
  2. 错误的更新方式:
    代码中每次都会创建新的MutableInteger替换map中的旧对象,而不是直接更新旧对象的value。即使MutableInteger的方法加了synchronized,也无法解决多线程替换对象导致的计数丢失问题。

  3. 额外隐患:SimpleDateFormat线程不安全:
    SimpleDateFormat不是线程安全的,多线程并发调用format方法会出现异常或格式化错误,虽然不是计数错误的直接原因,但会导致时间键生成错误,间接影响计数准确性。

修复方案

方案1:使用ConcurrentHashMap的原子更新方法(推荐)

利用ConcurrentHashMap提供的merge或compute方法,这两个方法保证整个更新流程的原子性,避免多线程覆盖问题。

修改后的saveFlow方法示例(改用线程安全的DateTimeFormatter替代SimpleDateFormat):

import java.time.LocalDateTime;
import java.time.ZoneId;
import java.time.format.DateTimeFormatter;
import java.util.Calendar;
import java.util.concurrent.ConcurrentHashMap;

// ...

// 接口流量记录
public static void saveFlow(Calendar c){
    System.out.println("Thread Name:"+Thread.currentThread().getName());
    // Java 8+ 线程安全的日期格式化器
    DateTimeFormatter dateFormatter = DateTimeFormatter.ofPattern("yyyy-MM-dd");
    DateTimeFormatter hourFormatter = DateTimeFormatter.ofPattern("yyyy-MM-dd HH");
    DateTimeFormatter minuteFormatter = DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm");

    LocalDateTime time = c.toInstant().atZone(ZoneId.systemDefault()).toLocalDateTime();
    String dayAcceptKey = dateFormatter.format(time);
    String hourAcceptKey = hourFormatter.format(time);
    String minuteAcceptKey = minuteFormatter.format(time);
    
    // merge方法原子更新:不存在则放入新值,存在则执行更新逻辑
    map.merge(minuteAcceptKey, new MutableInteger() {{ setValue(1); }}, (oldVal, newVal) -> {
        oldVal.setValue(oldVal.getValue() + 1);
        return oldVal;
    });
    System.out.println("minute in map:"+ minuteAcceptKey + ","+ map.get(minuteAcceptKey).getValue());
    
    map.merge(hourAcceptKey, new MutableInteger() {{ setValue(1); }}, (oldVal, newVal) -> {
        oldVal.setValue(oldVal.getValue() + 1);
        return oldVal;
    });
    System.out.println("hour in map:" + hourAcceptKey + "," + map.get(hourAcceptKey).getValue());
    
    map.merge(dayAcceptKey, new MutableInteger() {{ setValue(1); }}, (oldVal, newVal) -> {
        oldVal.setValue(oldVal.getValue() + 1);
        return oldVal;
    });
    System.out.println("day in map:" + dayAcceptKey + "," + map.get(dayAcceptKey).getValue());
}

方案2:替换为AtomicInteger

用Java内置的线程安全原子类AtomicInteger替代自定义的MutableInteger,结合merge方法更简洁:

// 修改map定义
static ConcurrentHashMap<String,AtomicInteger> map = new ConcurrentHashMap();

// ...

// saveFlow中的更新逻辑
map.merge(minuteAcceptKey, new AtomicInteger(1), (old, newOne) -> {
    old.incrementAndGet();
    return old;
});

方案3:修复SimpleDateFormat的线程安全问题

如果必须使用SimpleDateFormat,可以通过ThreadLocal为每个线程创建独立的实例:

private static ThreadLocal<SimpleDateFormat> dateFormatThreadLocal = ThreadLocal.withInitial(() -> new SimpleDateFormat("yyyy-MM-dd"));
private static ThreadLocal<SimpleDateFormat> hourFormatThreadLocal = ThreadLocal.withInitial(() -> new SimpleDateFormat("yyyy-MM-dd HH"));
private static ThreadLocal<SimpleDateFormat> minuteFormatThreadLocal = ThreadLocal.withInitial(() -> new SimpleDateFormat("yyyy-MM-dd HH:mm"));

// 使用时
String dayAcceptKey = dateFormatThreadLocal.get().format(c.getTime());

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.02 20:41:07