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()); } } }
问题原因分析
核心逻辑非原子性:
ConcurrentHashMap的put方法是原子操作,但整个计数流程是"创建新对象→put到map→获取旧对象→更新新对象值",这一系列操作不是原子的。多个线程同时处理同一个时间key时,会发生覆盖:- 线程A创建值为1的
MutableInteger,并put到map; - 线程B同时创建值为1的
MutableInteger,覆盖线程A的对象; - 线程A拿到旧值(可能是更早的对象),更新自己的新对象,但这个对象已经不在map里了,map里是线程B的对象,导致线程A的计数丢失。
- 线程A创建值为1的
错误的更新方式:
代码中每次都会创建新的MutableInteger替换map中的旧对象,而不是直接更新旧对象的value。即使MutableInteger的方法加了synchronized,也无法解决多线程替换对象导致的计数丢失问题。额外隐患: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

