에이전트와 비동기 액션

에이전트와 비동기 액션 (Agents and Asynchronous Actions)

Ref처럼 에이전트(agent)도 변경 가능한 상태에 대한 공유 접근을 제공해요. Ref여러 위치의 조정된(coordinated)·동기적 변경을 지원한다면, 에이전트는 개별 위치의 독립적·비동기적 변경을 제공해요. 에이전트는 수명 동안 단일 저장 위치에 바인딩되고, 그 위치의 변경(새 상태로)은 오직 액션(action)의 결과로만 일어나도록 허용해요. 액션은 (선택적으로 추가 인자를 가진) 함수로, 에이전트의 상태에 비동기적으로 적용되며 그 반환 값이 에이전트의 새 상태가 돼요. 액션은 함수이므로 멀티메서드일 수도 있고, 따라서 액션은 잠재적으로 다형(polymorphic)이에요. 또한 함수 집합이 열려 있으므로 에이전트가 지원하는 액션 집합도 열려 있어요. 이는 다른 언어들이 제공하는 패턴 매칭 메시지 처리 루프와 뚜렷한 대조를 이뤄요.

Clojure의 에이전트는 반응적(reactive) 이지 자율적(autonomous)이지 않아요 — 명령형 메시지 루프도, 블로킹 receive도 없어요. 에이전트의 상태는 그 자체로 불변이어야 하고(가급적 Clojure의 영속 컬렉션 중 하나의 인스턴스), 에이전트의 상태는 어떤 메시지 없이도 항상 어떤 스레드든 즉시 읽을 수 있어요(deref 함수나 리더 매크로 +@+ 사용) — 즉 관찰(observation)은 협력이나 조정이 필요 없어요.

에이전트 액션 디스패치는 +(send agent fn args*)+ 형태를 취해요. send(그리고 send-off)는 항상 즉시 반환해요. 조금 뒤, 다른 스레드에서 다음이 일어날 거예요:

. 주어진 +fn+이 에이전트의 _상태_와 제공된 인자(있다면)에 적용돼요. . +fn+의 반환 값이, 에이전트에 설정된 검증기 함수가 있다면 그 검증기 함수에 전달돼요. 자세한 내용은 set-validator!를 참고하세요. . 검증기가 성공하거나 검증기가 없으면, 주어진 +fn+의 반환 값이 에이전트의 새 상태가 돼요. . 에이전트에 추가된 감시자(watcher)가 있다면 호출돼요. 자세한 내용은 add-watch를 참고하세요. . 함수 실행 중에 (직접이든 간접이든) 다른 디스패치가 이루어지면, 에이전트의 상태가 변경된 이후까지 보류돼요.

액션 함수가 예외를 던지면 중첩 디스패치는 일어나지 않고, 예외가 에이전트 자체에 캐시돼요. 에이전트에 오류가 캐시되어 있으면, 그 후의 어떤 상호작용도 에이전트의 오류가 지워질 때까지 즉시 예외를 던져요. 에이전트 오류는 agent-error로 검사하고, 에이전트는 restart-agent로 재시작할 수 있어요.

모든 에이전트의 액션은 스레드 풀 안의 스레드들 사이에 인터리브돼요. 어느 시점에서든 각 에이전트에 대해 최대 하나의 액션이 실행되고 있어요. 한 에이전트나 스레드에서 에이전트로 디스패치된 액션들은 보낸 순서대로 발생하지만, 같은 에이전트에 다른 소스에서 디스패치된 액션들과 인터리브될 수 있어요. CPU에 제한된 액션에는 send를, IO에서 블로킹할 수 있는 액션에는 send-off가 적절해요.

에이전트는 STM과 통합돼요 — 트랜잭션 안에서 이루어진 디스패치는 커밋될 때까지 보류되고, 트랜잭션이 재시도되거나 중단되면 폐기돼요.

Clojure의 모든 동시성 지원과 마찬가지로, 사용자 코드 잠금(locking)은 관련되지 않아요.

에이전트를 사용하면 JVM 종료를 막는 non-daemon 백그라운드 스레드 풀이 시작된다는 점에 주의하세요. shutdown-agents를 사용해 이 스레드들을 종료하고 종료를 허용하세요.

출처: Clojure 공식 문서 - Agents and Asynchronous Actions

본문

예시 (Example)

이 예시는 링 주위로 메시지를 보내는(send-a-message-around-a-ring) 테스트의 구현이에요. m개의 에이전트 체인이 만들어지고, n개의 액션 시퀀스가 체인의 머리로 디스패치되어 그것을 통해 중계돼요.

(defn relay [x i]
  (when (:next x)
    (send (:next x) relay i))
  (when (and (zero? i) (:report-queue x))
    (.put (:report-queue x) i))
  x)

(defn run [m n]
  (let [q (new java.util.concurrent.SynchronousQueue)
        hd (reduce (fn [next _] (agent {:next next}))
                   (agent {:report-queue q}) (range (dec m)))]
    (doseq [i (reverse (range n))]
      (send hd relay i))
    (.take q)))

; 1 million message sends:
(time (run 1000 1000))
->"Elapsed time: 2959.254 msecs"

더 알아보기