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

如何用Java gRPC库在Clojure中消费SpiceDB LookupResources gRPC流

在Clojure中消费SpiceDB LookupResources gRPC流(无core.async)

同步阻塞实现(适合小结果集)

如果结果集较小可一次性加载,我们可以用Java同步工具CountDownLatch等待gRPC流结束,同时通过atom收集所有资源ID:

(ns your.namespace
  (:import [io.grpc.stub StreamObserver]
           [java.util.concurrent CountDownLatch]
           [com.authzed.api.v1 LookupResourcesRequest LookupResourcesResponse]))

(defn lookup-all-resources
  "调用SpiceDB的LookupResources gRPC端点,返回所有资源ID的向量"
  [spicedb-stub request]
  (let [results (atom [])
        latch (CountDownLatch. 1)
        error (atom nil)
        observer (reify StreamObserver
                   (onNext [_ ^LookupResourcesResponse resp]
                     ;; 从响应中提取资源ID并入结果集
                     (swap! results conj (.getResourceId resp)))
                   (onError [_ throwable]
                     (reset! error throwable)
                     (.countDown latch))
                   (onCompleted [_]
                     (.countDown latch)))]
    ;; 发起gRPC流请求
    (.lookupResources spicedb-stub request observer)
    ;; 等待流处理完成
    (.await latch)
    ;; 有错误则抛出
    (when-let [err @error]
      (throw (ex-info "gRPC流处理失败" {:error err})))
    ;; 返回结果向量
    @results))

懒加载实现(按需获取结果)

如果需要避免一次性加载大量数据,可通过LinkedBlockingQueue缓冲流结果,结合lazy-seq实现按需读取:

(ns your.namespace
  (:import [io.grpc.stub StreamObserver]
           [java.util.concurrent LinkedBlockingQueue]
           [com.authzed.api.v1 LookupResourcesRequest LookupResourcesResponse]))

;; 定义流结束标记
(def ^:private stream-end (Object.))

(defn lookup-resources-lazy
  "调用SpiceDB的LookupResources gRPC端点,返回懒加载的资源ID序列"
  [spicedb-stub request]
  (let [queue (LinkedBlockingQueue.)
        error (atom nil)
        observer (reify StreamObserver
                   (onNext [_ ^LookupResourcesResponse resp]
                     (.put queue (.getResourceId resp)))
                   (onError [_ throwable]
                     (reset! error throwable)
                     (.put queue stream-end))
                   (onCompleted [_]
                     (.put queue stream-end)))]
    ;; 异步发起gRPC流请求
    (.lookupResources spicedb-stub request observer)
    ;; 构建懒序列
    (lazy-seq
      (loop []
        (let [item (.take queue)]
          (cond
            (identical? item stream-end)
            (do (when-let [err @error]
                  (throw (ex-info "gRPC流处理失败" {:error err})))
                nil)
            :else
            (cons item (lazy-seq (loop)))))))))

关键说明

  • 两个实现均基于原生io.grpc库,无需引入core.async
  • 同步版本通过CountDownLatch阻塞主线程,直到gRPC流的onCompleted或onError回调触发
  • 懒加载版本利用阻塞队列和懒序列,每次消费时才取出下一个结果,适合处理大数据量
  • 两种实现都包含错误处理,gRPC流异常会转化为Clojure标准异常抛出

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.27 23:47:10