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

Clojure中ArrayBlockingQueue代理的多线程监控问题排查

Why Your ArrayBlockingQueue Proxy Fails in Multi-Threaded Scenarios

The core issue with your code is that ArrayBlockingQueue is a final class in Java. Clojure's proxy macro works by creating a subclass of the target class, but you can't subclass final classes—this violates Java's language rules and will throw a runtime exception (like IllegalArgumentException: Cannot subclass final class java.util.concurrent.ArrayBlockingQueue).

Even if that weren't the case, proxying concrete classes for method interception is often error-prone, especially with concurrent utilities that rely on internal locking. Instead, we should use the decorator pattern: implement the BlockingQueue interface, wrap an underlying ArrayBlockingQueue, and delegate all methods to it while intercepting the ones we want to monitor.

Corrected Code

Here's how to fix your implementation using Clojure's reify to implement the BlockingQueue interface:

(ns clj-super-bug.core
  (:import [java.util.concurrent ArrayBlockingQueue Executors BlockingQueue]))

(let [thread-count 10
      put-count 100
      executor (Executors/newFixedThreadPool thread-count)
      puts (atom 0)
      ;; Create the actual ArrayBlockingQueue we'll wrap
      inner-queue (ArrayBlockingQueue. 1000)
      ;; Implement BlockingQueue to decorate inner-queue
      queue (reify BlockingQueue
              ;; Intercept put to add monitoring
              (put [_ el]
                (.put inner-queue el)
                (swap! puts inc))
              ;; Delegate all other BlockingQueue methods to inner-queue
              (add [_ el] (.add inner-queue el))
              (addAll [_ coll] (.addAll inner-queue coll))
              (clear [_] (.clear inner-queue))
              (contains [_ o] (.contains inner-queue o))
              (containsAll [_ coll] (.containsAll inner-queue coll))
              (drainTo [_ c] (.drainTo inner-queue c))
              (drainTo [_ c maxElements] (.drainTo inner-queue c maxElements))
              (element [_] (.element inner-queue))
              (isEmpty [_] (.isEmpty inner-queue))
              (iterator [_] (.iterator inner-queue))
              (offer [_ el] (.offer inner-queue el))
              (offer [_ el timeout unit] (.offer inner-queue el timeout unit))
              (peek [_] (.peek inner-queue))
              (poll [_] (.poll inner-queue))
              (poll [_ timeout unit] (.poll inner-queue timeout unit))
              (remainingCapacity [_] (.remainingCapacity inner-queue))
              (remove [_ o] (.remove inner-queue o))
              (removeAll [_ coll] (.removeAll inner-queue coll))
              (retainAll [_ coll] (.retainAll inner-queue coll))
              (size [_] (.size inner-queue))
              (toArray [_] (.toArray inner-queue))
              (toArray [_ a] (.toArray inner-queue a))
              (take [_] (.take inner-queue)))]
  (.invokeAll executor (repeat put-count #(.put queue 0)))
  (assert (= (.size queue) put-count) "should have put in put-count items")
  (println "Success! Total puts:" @puts "Queue size:" (.size queue)))

Key Changes Explained

  1. Wrap Instead of Subclass: We create a real ArrayBlockingQueue instance (inner-queue) and use reify to implement the BlockingQueue interface around it. This avoids the final class restriction.
  2. Delegate All Methods: Every method of BlockingQueue is delegated to inner-queue except put, which we override to add our monitoring logic.
  3. Thread-Safe Monitoring: The puts atom is thread-safe, so incrementing it after a successful put works correctly even in multi-threaded scenarios.

This implementation will run without exceptions and pass your assertion, as all 100 elements will be added to the queue successfully.

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.25 04:21:50