From babc2ff7b87da64464d5d02ccd40314b15c7162c Mon Sep 17 00:00:00 2001 From: Magnus Landerblom Date: Wed, 8 Apr 2020 14:25:13 +0200 Subject: [PATCH 1/6] Add TCP connections For fetching metrics faster with simpler format --- src/cloudkarafka/fast_jmx.clj | 105 +++++++++++++++++++++++ src/cloudkarafka/kafka_http_reporter.clj | 26 ++++-- src/cloudkarafka/tcp.clj | 86 +++++++++++++++++++ 3 files changed, 210 insertions(+), 7 deletions(-) create mode 100644 src/cloudkarafka/fast_jmx.clj create mode 100644 src/cloudkarafka/tcp.clj diff --git a/src/cloudkarafka/fast_jmx.clj b/src/cloudkarafka/fast_jmx.clj new file mode 100644 index 0000000..58d2f56 --- /dev/null +++ b/src/cloudkarafka/fast_jmx.clj @@ -0,0 +1,105 @@ +(ns cloudkarafka.fast-jmx + (:require [clojure.java.jmx :as jmx] + [clojure.string :as str]) + (:import [java.io Reader Writer] + [javax.management ObjectName] + [java.util Map$Entry]) + ;(:gen-class) + ) + +(set! *warn-on-reflection* true) + +(defn write-err [str] + (let [^Writer o *err*] + (.write o (format "%s\n" str)) + (.flush o))) + +(defn parse-line [line] + (let [t (str/trim line) + [cmd bean] (str/split t #" ")] + [cmd bean])) + +(defn parse-host [str] + (let [[host port] (str/split str #":")] + {:host host + :port (Integer/parseInt port)})) + +(defn- mbean-keys->map [^Map$Entry e] + (vector (.getKey e) (.getValue e))) + +(defn valid-value? [v] + (or + (string? v) + (int? v) + (double? v) + (boolean? v))) + +(defn f + ([m] (f "" m)) + ([p m] + (reduce + (fn [res [k v]] + (if (map? v) + (merge res (f (str p (name k) ".") v)) + (if (valid-value? v) + (assoc res (str p (name k)) v) + res))) + {} + m))) + +(defn mbean-params [^ObjectName mbean] + (let [kpl (.getKeyPropertyList mbean) + entries (.entrySet kpl)] + (into {} (map mbean-keys->map entries)))) + +(defn bean-value [mbean] + (try + (jmx/mbean mbean) + (catch javax.management.OperationsException e + (write-err (.getMessage e)) + {}))) + +(defn jmx-values [mbean] + (println "REP" "beans..." mbean (jmx/mbean mbean)) + (let [params (mbean-params mbean) + values (f "" (jmx/mbean mbean))] + (println "REP" "res." params values) + (merge params values))) + +(defn query [^String bean] + (cond (.contains bean "*") (doall (map jmx-values (jmx/mbean-names bean))) + :else (list (jmx-values (jmx/as-object-name bean))))) + +(defn ^String format-result [values] + (->> values + (map (fn [[k v]] (str k "=" v))) + (str/join ";;"))) + +(defmulti exec (fn [a _] a)) + +(defmethod exec "jmx" [_ bean] + (query bean)) + +(defmethod exec "v" [_ bean] + (list {bean "2.3.1"})) + +(defmethod exec :default [_ _] + (list {})) + +(defn handler [^Reader in ^Writer out] + (doseq [ln (line-seq (java.io.BufferedReader. in))] + (do (println "REP" "line" ln) + (let [[cmd bean] (parse-line ln)] + (println "REP" "cmd" cmd bean) + (when-not (empty? cmd) + (println "REP" "run") + (doseq [res (exec cmd bean) + :when (seq res)] + (.write out (format-result res)) + (.write out "\n"))) + (.write out "\n") + (.flush out))))) + + + + diff --git a/src/cloudkarafka/kafka_http_reporter.clj b/src/cloudkarafka/kafka_http_reporter.clj index 9d1c7a6..5cc228d 100644 --- a/src/cloudkarafka/kafka_http_reporter.clj +++ b/src/cloudkarafka/kafka_http_reporter.clj @@ -6,7 +6,9 @@ [jsonista.core :as json] [compojure.core :refer :all] [compojure.route :as route] - [ring.middleware.params :as params]) + [ring.middleware.params :as params] + [cloudkarafka.tcp :as tcp] + [cloudkarafka.fast-jmx :as fjmx]) (:gen-class :implements [org.apache.kafka.common.metrics.MetricsReporter] :constructors {[] []})) @@ -75,6 +77,7 @@ (route/not-found "Not found")))) + (defn -configure [this config] (let [parsed-config (into {} (map (fn [[k v]] [(keyword k) v]) config)) uris (listener-uri (:listeners parsed-config) "PLAINTEXT") @@ -88,14 +91,23 @@ (defn -init [this metrics] (let [config (:kafka-config @state) - port (Integer/parseInt (or (:kafka_http_reporter.port config) "19092"))] - (println "[INFO] KafkaHttpReporter: Starting HTTP server on port " port ) - (swap! state assoc :http-server (http/start-server handler {:port port})))) + http-port (Integer/parseInt (or (:kafka_http_reporter.port config) "19092")) + tcp-port (Integer/parseInt (or (:kafka_http_reporter.tcp_port config) "19500")) + tcp-server (tcp/tcp-server :port tcp-port :handler (tcp/wrap-io fjmx/handler))] + (println "[INFO] KafkaHttpReporter: Starting HTTP server on port " http-port ) + (swap! state assoc :http-server (http/start-server handler {:port http-port})) + (println "[INFO] KafkaHttpReporter: Starting TCP server on port " tcp-port) + (swap! state assoc :tcp-server (tcp/start tcp-server)))) (defn -metricChange [this metric]) (defn -metricRemoval [this metric]) (defn -close [this] - (when-let [^java.io.Closeable s (:http-server @state)] - (println "[INFO] KafkaHttpReporter: Closing HTTP server") - (.close s))) + ;; No need to close this + #_(let [s @state] + (when-let [^java.io.Closeable server (:http-server s)] + (println "[INFO] KafkaHttpReporter: Closing HTTP server") + (.close server)) + (when-let [server (:tcp-server s)] + (println "[INFO] KafkaHttpReporter: Closing TCP server") + (tcp/stop server)))) diff --git a/src/cloudkarafka/tcp.clj b/src/cloudkarafka/tcp.clj new file mode 100644 index 0000000..72d43bf --- /dev/null +++ b/src/cloudkarafka/tcp.clj @@ -0,0 +1,86 @@ +(ns cloudkarafka.tcp + "Functions for creating a threaded TCP server." + (:require [clojure.java.io :as io]) + (:import [java.net InetAddress ServerSocket Socket SocketException])) + +(defn- server-socket [server] + (ServerSocket. + (:port server) + (:backlog server) + (InetAddress/getByName (:host server)))) + +(defn tcp-server + "Create a new TCP server. Takes the following keyword arguments: + :host - the host to bind to (defaults to 127.0.0.1) + :port - the port to bind to + :handler - a function to handle incoming connections, expects a socket as + an argument + :backlog - the maximum backlog of connections to keep (defaults to 50)" + [& {:as options}] + {:pre [(:port options) + (:handler options)]} + (merge + {:host "127.0.0.1" + :backlog 50 + :socket (atom nil) + :connections (atom #{})} + options)) + +(defn close-socket [server socket] + (swap! (:connections server) disj socket) + (when-not (.isClosed socket) + (.close socket))) + +(defn- open-server-socket [server] + (reset! (:socket server) + (server-socket server))) + +(defn- accept-connection + [{:keys [handler connections socket] :as server}] + (let [conn (.accept @socket)] + (swap! connections conj conn) + (future + (try (handler conn) + (finally (close-socket server conn)))))) + +(defn running? + "True if the server is running." + [server] + (if-let [socket @(:socket server)] + (not (.isClosed socket)))) + +(defn start + "Start a TCP server going." + [server] + (open-server-socket server) + (future + (while (running? server) + (try + (accept-connection server) + (catch SocketException _))))) + +(defn stop + "Stop the TCP server and close all open connections." + [server] + (doseq [socket @(:connections server)] + (close-socket server socket)) + (.close @(:socket server))) + +(defn wrap-streams + "Wrap a handler so that it expects an InputStream and an OutputStream + as arguments, rather than a raw Socket." + [handler] + (fn [socket] + (with-open [input (.getInputStream socket) + output (.getOutputStream socket)] + (handler input output)))) + +(defn wrap-io + "Wrap a handler so that it expects a Reader and Writer as arguments, rather + than a raw Socket." + [handler] + (wrap-streams + (fn [input output] + (with-open [reader (io/reader input) + writer (io/writer output)] + (handler reader writer))))) From 4400b31199bc72d6926918e0f5f43a86181867ce Mon Sep 17 00:00:00 2001 From: Magnus Landerblom Date: Tue, 14 Apr 2020 21:51:22 +0200 Subject: [PATCH 2/6] Changes --- src/cloudkarafka/fast_jmx.clj | 37 ++++++++++++----------------------- 1 file changed, 13 insertions(+), 24 deletions(-) diff --git a/src/cloudkarafka/fast_jmx.clj b/src/cloudkarafka/fast_jmx.clj index 58d2f56..4a8367e 100644 --- a/src/cloudkarafka/fast_jmx.clj +++ b/src/cloudkarafka/fast_jmx.clj @@ -3,9 +3,7 @@ [clojure.string :as str]) (:import [java.io Reader Writer] [javax.management ObjectName] - [java.util Map$Entry]) - ;(:gen-class) - ) + [java.util Map$Entry])) (set! *warn-on-reflection* true) @@ -34,13 +32,13 @@ (double? v) (boolean? v))) -(defn f - ([m] (f "" m)) +(defn build-value-map + ([m] (build-value-map "" m)) ([p m] (reduce (fn [res [k v]] (if (map? v) - (merge res (f (str p (name k) ".") v)) + (merge res (build-value-map (str p (name k) ".") v)) (if (valid-value? v) (assoc res (str p (name k)) v) res))) @@ -60,10 +58,8 @@ {}))) (defn jmx-values [mbean] - (println "REP" "beans..." mbean (jmx/mbean mbean)) (let [params (mbean-params mbean) - values (f "" (jmx/mbean mbean))] - (println "REP" "res." params values) + values (build-value-map "" (jmx/mbean mbean))] (merge params values))) (defn query [^String bean] @@ -88,18 +84,11 @@ (defn handler [^Reader in ^Writer out] (doseq [ln (line-seq (java.io.BufferedReader. in))] - (do (println "REP" "line" ln) - (let [[cmd bean] (parse-line ln)] - (println "REP" "cmd" cmd bean) - (when-not (empty? cmd) - (println "REP" "run") - (doseq [res (exec cmd bean) - :when (seq res)] - (.write out (format-result res)) - (.write out "\n"))) - (.write out "\n") - (.flush out))))) - - - - + (let [[cmd bean] (parse-line ln)] + (when-not (empty? cmd) + (doseq [res (exec cmd bean) + :when (seq res)] + (.write out (format-result res)) + (.write out "\n"))) + (.write out "\n") + (.flush out)))) From a3650001df1b04a42e4ebc4bb97dc8180b4efcbb Mon Sep 17 00:00:00 2001 From: Magnus Landerblom Date: Wed, 15 Apr 2020 23:08:47 +0200 Subject: [PATCH 3/6] Use aleph.tcp instead --- project.clj | 7 +- src/cloudkarafka/{fast_jmx.clj => cmds.clj} | 42 +++--- src/cloudkarafka/kafka_http_reporter.clj | 81 +++++------- src/cloudkarafka/kafkaadmin.clj | 126 +++++++++++------- src/cloudkarafka/tcp.clj | 139 +++++++++----------- src/cloudkarafka/util.clj | 20 +++ 6 files changed, 211 insertions(+), 204 deletions(-) rename src/cloudkarafka/{fast_jmx.clj => cmds.clj} (66%) create mode 100644 src/cloudkarafka/util.clj diff --git a/project.clj b/project.clj index 7b91a95..a669fbd 100644 --- a/project.clj +++ b/project.clj @@ -6,9 +6,10 @@ :profiles {:uberjar {:aot :all} :provided {:dependencies [[org.apache.kafka/kafka_2.11 "1.1.0"] [org.apache.kafka/kafka-clients "1.1.0"]]}} - :dependencies [[org.clojure/clojure "1.9.0"] - [org.clojure/java.jmx "0.3.4"] - [metosin/jsonista "0.2.2"] + :dependencies [[org.clojure/clojure "1.10.1"] + [org.clojure/java.jmx "1.0.0"] + [metosin/jsonista "0.2.5"] [aleph "0.4.6"] + [manifold "0.1.8"] [compojure "1.6.1"] [t6/from-scala "0.3.0"]]) diff --git a/src/cloudkarafka/fast_jmx.clj b/src/cloudkarafka/cmds.clj similarity index 66% rename from src/cloudkarafka/fast_jmx.clj rename to src/cloudkarafka/cmds.clj index 4a8367e..726abe3 100644 --- a/src/cloudkarafka/fast_jmx.clj +++ b/src/cloudkarafka/cmds.clj @@ -1,27 +1,16 @@ -(ns cloudkarafka.fast-jmx +(ns cloudkarafka.cmds (:require [clojure.java.jmx :as jmx] + [cloudkarafka.kafkaadmin :as ka] [clojure.string :as str]) - (:import [java.io Reader Writer] + (:import java.io.Writer [javax.management ObjectName] [java.util Map$Entry])) -(set! *warn-on-reflection* true) - (defn write-err [str] (let [^Writer o *err*] (.write o (format "%s\n" str)) (.flush o))) -(defn parse-line [line] - (let [t (str/trim line) - [cmd bean] (str/split t #" ")] - [cmd bean])) - -(defn parse-host [str] - (let [[host port] (str/split str #":")] - {:host host - :port (Integer/parseInt port)})) - (defn- mbean-keys->map [^Map$Entry e] (vector (.getKey e) (.getValue e))) @@ -66,29 +55,28 @@ (cond (.contains bean "*") (doall (map jmx-values (jmx/mbean-names bean))) :else (list (jmx-values (jmx/as-object-name bean))))) -(defn ^String format-result [values] +(defn format-result-row [values] (->> values - (map (fn [[k v]] (str k "=" v))) + (map (fn [[k v]] (str (name k) "=" v))) (str/join ";;"))) +(defn format-result [values] + (str/join "\n" (map format-result-row values))) + (defmulti exec (fn [a _] a)) (defmethod exec "jmx" [_ bean] (query bean)) (defmethod exec "v" [_ bean] - (list {bean "2.3.1"})) + (list (case bean + "kafka" {"kafka" (org.apache.kafka.common.utils.AppInfoParser/getVersion)} + "plugin" {"plugin" "1.0.0"}))) + +(defmethod exec "groups" [_ group] + (ka/consumers)) (defmethod exec :default [_ _] (list {})) -(defn handler [^Reader in ^Writer out] - (doseq [ln (line-seq (java.io.BufferedReader. in))] - (let [[cmd bean] (parse-line ln)] - (when-not (empty? cmd) - (doseq [res (exec cmd bean) - :when (seq res)] - (.write out (format-result res)) - (.write out "\n"))) - (.write out "\n") - (.flush out)))) + diff --git a/src/cloudkarafka/kafka_http_reporter.clj b/src/cloudkarafka/kafka_http_reporter.clj index 5cc228d..955944b 100644 --- a/src/cloudkarafka/kafka_http_reporter.clj +++ b/src/cloudkarafka/kafka_http_reporter.clj @@ -1,50 +1,22 @@ (ns cloudkarafka.kafka-http-reporter (:require [cloudkarafka.jmx :as jmx] [cloudkarafka.kafkaadmin :as ka] + [cloudkarafka.cmds :refer [exec format-result]] + [cloudkarafka.util :as util] [clojure.string :as str] [aleph.http :as http] + [aleph.tcp :as tcp] + [manifold.stream :as s] [jsonista.core :as json] [compojure.core :refer :all] [compojure.route :as route] - [ring.middleware.params :as params] - [cloudkarafka.tcp :as tcp] - [cloudkarafka.fast-jmx :as fjmx]) + [ring.middleware.params :as params]) (:gen-class :implements [org.apache.kafka.common.metrics.MetricsReporter] :constructors {[] []})) (set! *warn-on-reflection* true) -(def mapper (json/object-mapper {:decode-key-fn true, :encode-key-fn true})) - -(def state (atom nil)) - -(defn listener-uri [listeners type] - (->> (str/split listeners #",") - (map #(str/split % #"://")) - (filter #(= (first %) type)) - (map second))) - -(defn modern-kafka? [version] - (let [major (-> version - (str/split #"\.") - first - Integer/parseInt)] - (>= major 2))) - -(defn consumers - ([] (consumers :group)) - ([group-by-fn] - (let [s @state - member-list (if (modern-kafka? (:kafka-version s)) - (ka/consumer-groups (:admin-client s) (:consumer s)) - (let [plaintext-url (-> (:kafka-config s) - :listeners - (listener-uri "PLAINTEXT") - first)] - (ka/consumer-groups-old plaintext-url (:consumer s))))] - (group-by group-by-fn member-list)))) - (def handler (params/wrap-params (routes @@ -52,57 +24,66 @@ {:status 200 :headers {"content-type" "text/plain"} :body "0.1.0"}) (GET "/kafka-version" [] - {:status 200 :headers {"content-type" "text/plain"} :body (:kafka-version @state)}) + {:status 200 :headers {"content-type" "text/plain"} :body (:kafka-version @util/state)}) (GET "/jmx" [bean group attrs] (if-let [values (jmx/query bean group (str/split attrs #","))] {:status 200 :headers {"content-type" "application/json"} - :body (json/write-value-as-string values mapper)} + :body (json/write-value-as-string values util/mapper)} {:status 404 :body (str "Bean " bean " not found")})) (GET "/config" [] - (if-let [c (:kafka-config @state)] + (if-let [c (:kafka-config @util/state)] {:status 200 :headers {"content-type" "application/json"} - :body (json/write-value-as-string c mapper)} + :body (json/write-value-as-string c util/mapper)} {:status 404 :body "No config"})) (GET "/consumer-groups" [] {:status 200 :headers {"content-type" "application/json"} - :body (json/write-value-as-string (consumers) mapper)}) + :body (json/write-value-as-string (group-by :group (ka/consumers)) util/mapper)}) (route/not-found "Not found")))) +(defn wrapper [f] + (fn [s _info] + (s/connect (s/map f s) s))) + +(defn tcp-handler [^bytes in] + (let [s (str/trim (String. in)) + [cmd bean] (str/split s #" ")] + (when-not (empty? cmd) + (str (format-result (exec cmd bean)) "\n\n")))) -(defn -configure [this config] +(defn -configure [_ config] (let [parsed-config (into {} (map (fn [[k v]] [(keyword k) v]) config)) - uris (listener-uri (:listeners parsed-config) "PLAINTEXT") + uris (util/listener-uri (:listeners parsed-config) "PLAINTEXT") props {:bootstrap.servers (first uris)} kafka-version (org.apache.kafka.common.utils.AppInfoParser/getVersion)] - (reset! state {:kafka-version kafka-version + (reset! util/state {:kafka-version kafka-version :kafka-config parsed-config - :admin-client (when (modern-kafka? kafka-version) + :admin-client (when (util/modern-kafka? kafka-version) (ka/admin-client props)) :consumer (ka/kafka-consumer props)}))) -(defn -init [this metrics] - (let [config (:kafka-config @state) +(defn -init [_ _] + (let [config (:kafka-config @util/state) http-port (Integer/parseInt (or (:kafka_http_reporter.port config) "19092")) tcp-port (Integer/parseInt (or (:kafka_http_reporter.tcp_port config) "19500")) - tcp-server (tcp/tcp-server :port tcp-port :handler (tcp/wrap-io fjmx/handler))] + tcp-server (tcp/start-server (wrapper tcp-handler) {:port tcp-port})] (println "[INFO] KafkaHttpReporter: Starting HTTP server on port " http-port ) - (swap! state assoc :http-server (http/start-server handler {:port http-port})) + (swap! util/state assoc :http-server (http/start-server handler {:port http-port})) (println "[INFO] KafkaHttpReporter: Starting TCP server on port " tcp-port) - (swap! state assoc :tcp-server (tcp/start tcp-server)))) + (swap! util/state assoc :tcp-server tcp-server))) -(defn -metricChange [this metric]) -(defn -metricRemoval [this metric]) +(defn -metricChange [_ _]) +(defn -metricRemoval [_ _]) -(defn -close [this] +(defn -close [_] ;; No need to close this #_(let [s @state] (when-let [^java.io.Closeable server (:http-server s)] diff --git a/src/cloudkarafka/kafkaadmin.clj b/src/cloudkarafka/kafkaadmin.clj index f02d13e..047916b 100644 --- a/src/cloudkarafka/kafkaadmin.clj +++ b/src/cloudkarafka/kafkaadmin.clj @@ -1,5 +1,6 @@ (ns cloudkarafka.kafkaadmin - (:require [t6.from-scala.core :as $])) + (:require [cloudkarafka.util :as util] + [t6.from-scala.core :as $])) (defn map->props [m] @@ -23,54 +24,59 @@ :consumerid (.consumerId member), :host (.host member)} (for [tp toppar - :let [group-partitions (get group-offset tp) - go (when group-partitions (.offset group-partitions)) - lo (get log-offset tp)]] - {:group group-name - :topic (.topic tp) - :partition (.partition tp) - :current_offset go, - :log_end_offset lo, - :lag (when go (- lo go)), - :clientid (.clientId member), - :consumerid (.consumerId member), - :host (.host member)})))) + :let [group-partitions (get group-offset tp) + go (when group-partitions (.offset group-partitions)) + lo (get log-offset tp)]] + {:group group-name + :topic (.topic tp) + :partition (.partition tp) + :current_offset go, + :log_end_offset lo, + :lag (when go (- lo go)), + :clientid (.clientId member), + :consumerid (.consumerId member), + :host (.host member)})))) + (defn consumer-groups - [client consumer] - (let [group-ids (into [] (map #(.groupId %) (.get (.all (.listConsumerGroups client))))) - descs (.get (.all (.describeConsumerGroups client group-ids)))] - (flatten (for [group-id group-ids - :let [desc (get descs group-id) - group-offset (.get (.partitionsToOffsetAndMetadata (.listConsumerGroupOffsets client group-id))) - log-offset (.endOffsets consumer (mapv key group-offset))]] - (member-list group-id desc group-offset log-offset))))) + ([client consumer] + (consumer-groups + client + consumer + (into [] (map #(.groupId %) (.get (.all (.listConsumerGroups client))))))) + ([client consumer group-ids] + (let [descs (.get (.all (.describeConsumerGroups client group-ids)))] + (flatten (for [group-id group-ids + :let [desc (get descs group-id) + group-offset (.get (.partitionsToOffsetAndMetadata (.listConsumerGroupOffsets client group-id))) + log-offset (.endOffsets consumer (mapv key group-offset))]] + (member-list group-id desc group-offset log-offset)))))) (defn consumer-groups-old - [url consumer] - (with-open [client (kafka.admin.AdminClient/createSimplePlaintext url)] - (let [res (java.util.LinkedList.)] - ($/for [group (.listAllConsumerGroupsFlattened client) - :let [summary (.describeConsumerGroup client (.groupId group) 0) - group-offset (.listGroupOffsets client (.groupId group))]] - ($/if-let [members (.consumers summary)] - ($/for [member members - :let [log-end-offsets (.endOffsets consumer (scala.collection.JavaConversions/asJavaCollection (.assignment member)))]] - ($/for [toppar (.assignment member) - :let [current-offset (get (scala.collection.JavaConversions/mapAsJavaMap group-offset) toppar) - log-end (get log-end-offsets toppar)]] - (.add res {:state (.state summary) - :group (.groupId group) - :topic (.topic toppar) - :partition (.partition toppar) - :current_offset current-offset - :log_end_offset log-end - :lag (and log-end current-offset (- log-end current-offset)) - :consumerid (.consumerId member) - :clientid (.clientId member) - :host (.host member)}))) - :none)) - res))) + ([url consumer] + (with-open [client (kafka.admin.AdminClient/createSimplePlaintext url)] + (let [res (java.util.LinkedList.)] + ($/for [group (.listAllConsumerGroupsFlattened client) + :let [summary (.describeConsumerGroup client (.groupId group) 0) + group-offset (.listGroupOffsets client (.groupId group))]] + ($/if-let [members (.consumers summary)] + ($/for [member members + :let [log-end-offsets (.endOffsets consumer (scala.collection.JavaConversions/asJavaCollection (.assignment member)))]] + ($/for [toppar (.assignment member) + :let [current-offset (get (scala.collection.JavaConversions/mapAsJavaMap group-offset) toppar) + log-end (get log-end-offsets toppar)]] + (.add res {:state (.state summary) + :group (.groupId group) + :topic (.topic toppar) + :partition (.partition toppar) + :current_offset current-offset + :log_end_offset log-end + :lag (and log-end current-offset (- log-end current-offset)) + :consumerid (.consumerId member) + :clientid (.clientId member) + :host (.host member)}))) + :none)) + res)))) (defn admin-client ^org.apache.kafka.clients.admin.AdminClient [props] @@ -85,8 +91,32 @@ props (map->props m)] (org.apache.kafka.clients.consumer.KafkaConsumer. props))) -(comment +(defn consumers + ([] + (let [s @util/state] + (if (util/modern-kafka? (:kafka-version s)) + (consumer-groups (:admin-client s) (:consumer s)) + (let [plaintext-url (-> (:kafka-config s) + :listeners + (util/listener-uri "PLAINTEXT") + first)] + (consumer-groups-old plaintext-url (:consumer s)))))) + ([group] + (let [s @util/state] + (if (util/modern-kafka? (:kafka-version s)) + (consumer-group (:admin-client s) (:consumer s) group) + (let [plaintext-url (-> (:kafka-config s) + :listeners + (util/listener-uri "PLAINTEXT") + first)] + (consumer-group-old plaintext-url (:consumer s) group)))))) - (consumer-groups-old "127.0.0.1:9092" (kafka-consumer {:bootstrap.servers "127.0.0.1:9092"})) +(comment + (consumer-groups-old + "127.0.0.1:9092" + (kafka-consumer {:bootstrap.servers "127.0.0.1:9092"})) - ) + (consumer-groups + (admin-client {:bootstrap.servers "127.0.0.1:9092"}) + (kafka-consumer {:bootstrap.servers "127.0.0.1:9092"})) +) diff --git a/src/cloudkarafka/tcp.clj b/src/cloudkarafka/tcp.clj index 72d43bf..7dc054b 100644 --- a/src/cloudkarafka/tcp.clj +++ b/src/cloudkarafka/tcp.clj @@ -1,86 +1,73 @@ -(ns cloudkarafka.tcp - "Functions for creating a threaded TCP server." - (:require [clojure.java.io :as io]) - (:import [java.net InetAddress ServerSocket Socket SocketException])) +(ns cloudkarafka.cmds + (:require [clojure.java.jmx :as jmx] + [cloudkarafka.kafkaadmin :as ka] + [clojure.string :as str]) + (:import java.io.Writer + [javax.management ObjectName] + [java.util Map$Entry]) + (:gen-class)) -(defn- server-socket [server] - (ServerSocket. - (:port server) - (:backlog server) - (InetAddress/getByName (:host server)))) +(defn write-err [str] + (let [^Writer o *err*] + (.write o (format "%s\n" str)) + (.flush o))) -(defn tcp-server - "Create a new TCP server. Takes the following keyword arguments: - :host - the host to bind to (defaults to 127.0.0.1) - :port - the port to bind to - :handler - a function to handle incoming connections, expects a socket as - an argument - :backlog - the maximum backlog of connections to keep (defaults to 50)" - [& {:as options}] - {:pre [(:port options) - (:handler options)]} - (merge - {:host "127.0.0.1" - :backlog 50 - :socket (atom nil) - :connections (atom #{})} - options)) +(defn- mbean-keys->map [^Map$Entry e] + (vector (.getKey e) (.getValue e))) -(defn close-socket [server socket] - (swap! (:connections server) disj socket) - (when-not (.isClosed socket) - (.close socket))) +(defn valid-value? [v] + (or + (string? v) + (int? v) + (double? v) + (boolean? v))) -(defn- open-server-socket [server] - (reset! (:socket server) - (server-socket server))) +(defn build-value-map + ([m] (build-value-map "" m)) + ([p m] + (reduce + (fn [res [k v]] + (if (map? v) + (merge res (build-value-map (str p (name k) ".") v)) + (if (valid-value? v) + (assoc res (str p (name k)) v) + res))) + {} + m))) -(defn- accept-connection - [{:keys [handler connections socket] :as server}] - (let [conn (.accept @socket)] - (swap! connections conj conn) - (future - (try (handler conn) - (finally (close-socket server conn)))))) +(defn mbean-params [^ObjectName mbean] + (let [kpl (.getKeyPropertyList mbean) + entries (.entrySet kpl)] + (into {} (map mbean-keys->map entries)))) -(defn running? - "True if the server is running." - [server] - (if-let [socket @(:socket server)] - (not (.isClosed socket)))) +(defn bean-value [mbean] + (try + (jmx/mbean mbean) + (catch javax.management.OperationsException e + (write-err (.getMessage e)) + {}))) -(defn start - "Start a TCP server going." - [server] - (open-server-socket server) - (future - (while (running? server) - (try - (accept-connection server) - (catch SocketException _))))) +(defn jmx-values [mbean] + (let [params (mbean-params mbean) + values (build-value-map "" (jmx/mbean mbean))] + (merge params values))) -(defn stop - "Stop the TCP server and close all open connections." - [server] - (doseq [socket @(:connections server)] - (close-socket server socket)) - (.close @(:socket server))) +(defn query [^String bean] + (cond (.contains bean "*") (doall (map jmx-values (jmx/mbean-names bean))) + :else (list (jmx-values (jmx/as-object-name bean))))) + +(defmulti exec (fn [a _] a)) + +(defmethod exec "jmx" [_ bean] + (query bean)) + +(defmethod exec "version" [_ bean] + (list {bean "2.3.1"})) + +(defmethod exec "groups" [_ bean] + (ka/consumers)) + +(defmethod exec :default [_ _] + (list {})) -(defn wrap-streams - "Wrap a handler so that it expects an InputStream and an OutputStream - as arguments, rather than a raw Socket." - [handler] - (fn [socket] - (with-open [input (.getInputStream socket) - output (.getOutputStream socket)] - (handler input output)))) -(defn wrap-io - "Wrap a handler so that it expects a Reader and Writer as arguments, rather - than a raw Socket." - [handler] - (wrap-streams - (fn [input output] - (with-open [reader (io/reader input) - writer (io/writer output)] - (handler reader writer))))) diff --git a/src/cloudkarafka/util.clj b/src/cloudkarafka/util.clj new file mode 100644 index 0000000..81517d5 --- /dev/null +++ b/src/cloudkarafka/util.clj @@ -0,0 +1,20 @@ +(ns cloudkarafka.util + (:require [jsonista.core :as json] + [clojure.string :as str])) + +(def mapper (json/object-mapper {:decode-key-fn true, :encode-key-fn true})) + +(def state (atom nil)) + +(defn listener-uri [listeners type] + (->> (str/split listeners #",") + (map #(str/split % #"://")) + (filter #(= (first %) type)) + (map second))) + +(defn modern-kafka? [version] + (let [major (-> version + (str/split #"\.") + first + Integer/parseInt)] + (>= major 2))) From ab52787ac7a71ae95ee28d075c39364225a527bf Mon Sep 17 00:00:00 2001 From: Magnus Landerblom Date: Fri, 17 Apr 2020 08:58:11 +0200 Subject: [PATCH 4/6] Consumer groups over TCP --- project.clj | 4 +- src/cloudkarafka/cmds.clj | 35 +++++++---------- src/cloudkarafka/kafka_http_reporter.clj | 30 +++++++-------- src/cloudkarafka/kafkaadmin.clj | 49 ++++++++++++------------ 4 files changed, 54 insertions(+), 64 deletions(-) diff --git a/project.clj b/project.clj index a669fbd..e102b67 100644 --- a/project.clj +++ b/project.clj @@ -1,6 +1,6 @@ -(defproject kafka-http-reporter "0.3.3" +(defproject kafka-http-reporter "0.4.0" :description "Expose JMX metrics through a HTTP interface" - :url "http://github.com/CloudKarafka/kafka-http-reporter" + :url "http://github.com/CloudKarafka/KafkaHttpReporter" :license {:name "Apache License 2.0" :url "https://github.com/CloudKarafka/KafkaHttpReporter/blob/master/LICENSE"} :profiles {:uberjar {:aot :all} diff --git a/src/cloudkarafka/cmds.clj b/src/cloudkarafka/cmds.clj index 726abe3..479ff8d 100644 --- a/src/cloudkarafka/cmds.clj +++ b/src/cloudkarafka/cmds.clj @@ -2,15 +2,9 @@ (:require [clojure.java.jmx :as jmx] [cloudkarafka.kafkaadmin :as ka] [clojure.string :as str]) - (:import java.io.Writer - [javax.management ObjectName] + (:import [javax.management ObjectName] [java.util Map$Entry])) -(defn write-err [str] - (let [^Writer o *err*] - (.write o (format "%s\n" str)) - (.flush o))) - (defn- mbean-keys->map [^Map$Entry e] (vector (.getKey e) (.getValue e))) @@ -39,13 +33,6 @@ entries (.entrySet kpl)] (into {} (map mbean-keys->map entries)))) -(defn bean-value [mbean] - (try - (jmx/mbean mbean) - (catch javax.management.OperationsException e - (write-err (.getMessage e)) - {}))) - (defn jmx-values [mbean] (let [params (mbean-params mbean) values (build-value-map "" (jmx/mbean mbean))] @@ -61,22 +48,28 @@ (str/join ";;"))) (defn format-result [values] - (str/join "\n" (map format-result-row values))) + (str + (str/join "\n" (map format-result-row values)) + "\n")) (defmulti exec (fn [a _] a)) (defmethod exec "jmx" [_ bean] - (query bean)) + (try + (when bean + (query bean)) + (catch javax.management.OperationsException e + (println e)))) (defmethod exec "v" [_ bean] - (list (case bean - "kafka" {"kafka" (org.apache.kafka.common.utils.AppInfoParser/getVersion)} - "plugin" {"plugin" "1.0.0"}))) + (case bean + "kafka" (list {"kafka" (org.apache.kafka.common.utils.AppInfoParser/getVersion)}) + :else nil)) (defmethod exec "groups" [_ group] - (ka/consumers)) + (ka/consumers group)) (defmethod exec :default [_ _] - (list {})) + nil) diff --git a/src/cloudkarafka/kafka_http_reporter.clj b/src/cloudkarafka/kafka_http_reporter.clj index 955944b..5c53077 100644 --- a/src/cloudkarafka/kafka_http_reporter.clj +++ b/src/cloudkarafka/kafka_http_reporter.clj @@ -57,11 +57,14 @@ (let [s (str/trim (String. in)) [cmd bean] (str/split s #" ")] (when-not (empty? cmd) - (str (format-result (exec cmd bean)) "\n\n")))) + (let [res (exec cmd bean)] + (str (when res (format-result res)) "\n"))))) (defn -configure [_ config] + (println "[INFO] KafkaHttpReporter: configure") (let [parsed-config (into {} (map (fn [[k v]] [(keyword k) v]) config)) - uris (util/listener-uri (:listeners parsed-config) "PLAINTEXT") + listener-name (or (:security.inter.broker.protocol parsed-config) "PLAINTEXT") + uris (util/listener-uri (:listeners parsed-config) listener-name) props {:bootstrap.servers (first uris)} kafka-version (org.apache.kafka.common.utils.AppInfoParser/getVersion)] (reset! util/state {:kafka-version kafka-version @@ -72,23 +75,18 @@ (defn -init [_ _] (let [config (:kafka-config @util/state) - http-port (Integer/parseInt (or (:kafka_http_reporter.port config) "19092")) - tcp-port (Integer/parseInt (or (:kafka_http_reporter.tcp_port config) "19500")) - tcp-server (tcp/start-server (wrapper tcp-handler) {:port tcp-port})] + http-port (Integer/parseInt (or (:kafka_http_reporter.port config) + (:kafkahttpreporter.port config) + "19092")) + tcp-port (Integer/parseInt (or (:kafka_http_reporter.tcp_port config) + (:kafkahttpreporter.tcp_port config) + "19500"))] (println "[INFO] KafkaHttpReporter: Starting HTTP server on port " http-port ) (swap! util/state assoc :http-server (http/start-server handler {:port http-port})) + (println "[INFO] KafkaHttpReporter: Starting TCP server on port " tcp-port) - (swap! util/state assoc :tcp-server tcp-server))) + (swap! util/state assoc :tcp-server (tcp/start-server (wrapper tcp-handler) {:port tcp-port})))) (defn -metricChange [_ _]) (defn -metricRemoval [_ _]) - -(defn -close [_] - ;; No need to close this - #_(let [s @state] - (when-let [^java.io.Closeable server (:http-server s)] - (println "[INFO] KafkaHttpReporter: Closing HTTP server") - (.close server)) - (when-let [server (:tcp-server s)] - (println "[INFO] KafkaHttpReporter: Closing TCP server") - (tcp/stop server)))) +(defn -close [_]) diff --git a/src/cloudkarafka/kafkaadmin.clj b/src/cloudkarafka/kafkaadmin.clj index 047916b..e5c5154 100644 --- a/src/cloudkarafka/kafkaadmin.clj +++ b/src/cloudkarafka/kafkaadmin.clj @@ -37,28 +37,33 @@ :consumerid (.consumerId member), :host (.host member)})))) - (defn consumer-groups ([client consumer] - (consumer-groups - client - consumer - (into [] (map #(.groupId %) (.get (.all (.listConsumerGroups client))))))) + (consumer-groups client consumer nil)) ([client consumer group-ids] - (let [descs (.get (.all (.describeConsumerGroups client group-ids)))] + (let [group-ids (or group-ids + (into [] (map #(.groupId %) (.get (.all (.listConsumerGroups client)))))) + descs (.get (.all (.describeConsumerGroups client group-ids)))] (flatten (for [group-id group-ids :let [desc (get descs group-id) group-offset (.get (.partitionsToOffsetAndMetadata (.listConsumerGroupOffsets client group-id))) log-offset (.endOffsets consumer (mapv key group-offset))]] (member-list group-id desc group-offset log-offset)))))) +(defn filtered-groups [client group-ids] + (let [all-group-ids (map #(.groupId %) (.listAllConsumerGroupsFlattened client)) + wanted (set group-ids)] + (if (empty? wanted) + all-group-ids + (filter #(contains? wanted %) all-group-ids)))) + (defn consumer-groups-old - ([url consumer] + ([url consumer group-ids] (with-open [client (kafka.admin.AdminClient/createSimplePlaintext url)] (let [res (java.util.LinkedList.)] - ($/for [group (.listAllConsumerGroupsFlattened client) - :let [summary (.describeConsumerGroup client (.groupId group) 0) - group-offset (.listGroupOffsets client (.groupId group))]] + ($/for [group-id (filtered-groups client group-ids) + :let [summary (.describeConsumerGroup client group-id 0) + group-offset (.listGroupOffsets client group-id)]] ($/if-let [members (.consumers summary)] ($/for [member members :let [log-end-offsets (.endOffsets consumer (scala.collection.JavaConversions/asJavaCollection (.assignment member)))]] @@ -66,7 +71,7 @@ :let [current-offset (get (scala.collection.JavaConversions/mapAsJavaMap group-offset) toppar) log-end (get log-end-offsets toppar)]] (.add res {:state (.state summary) - :group (.groupId group) + :group group-id :topic (.topic toppar) :partition (.partition toppar) :current_offset current-offset @@ -92,24 +97,18 @@ (org.apache.kafka.clients.consumer.KafkaConsumer. props))) (defn consumers - ([] - (let [s @util/state] - (if (util/modern-kafka? (:kafka-version s)) - (consumer-groups (:admin-client s) (:consumer s)) - (let [plaintext-url (-> (:kafka-config s) - :listeners - (util/listener-uri "PLAINTEXT") - first)] - (consumer-groups-old plaintext-url (:consumer s)))))) + ([] (consumers nil)) ([group] - (let [s @util/state] + (let [s @util/state + kafka-config (:kafka-config s)] (if (util/modern-kafka? (:kafka-version s)) - (consumer-group (:admin-client s) (:consumer s) group) - (let [plaintext-url (-> (:kafka-config s) + (consumer-groups (:admin-client s) (:consumer s) group) + (let [plaintext-url (-> kafka-config :listeners - (util/listener-uri "PLAINTEXT") + (util/listener-uri + (or (:security.inter.broker.protocol kafka-config) "PLAINTEXT")) first)] - (consumer-group-old plaintext-url (:consumer s) group)))))) + (consumer-groups-old plaintext-url (:consumer s) group)))))) (comment (consumer-groups-old From b33f16eddcc54205d8daf3c9186f1832d819fbea Mon Sep 17 00:00:00 2001 From: Magnus Landerblom Date: Fri, 17 Apr 2020 13:48:28 +0200 Subject: [PATCH 5/6] work --- README.md | 52 +++++++++++++++------ src/cloudkarafka/cmds.clj | 13 ++++-- src/cloudkarafka/kafka_http_reporter.clj | 22 ++++----- src/cloudkarafka/kafkaadmin.clj | 59 ++++++------------------ src/cloudkarafka/util.clj | 16 ++++--- 5 files changed, 81 insertions(+), 81 deletions(-) diff --git a/README.md b/README.md index b4e077d..56e73b0 100644 --- a/README.md +++ b/README.md @@ -2,32 +2,56 @@ A metrics reporter to be installed on your kafka brokers -It exposes an HTTP API that is used by https://github.com/CloudKarafka/cloudkarafka-manager to get metrics and broker information. +It exposes an HTTP API and a TCP server that is used by for example https://github.com/CloudKarafka/cloudkarafka-manager to get metrics and broker information. + +## TCP actions + +``` +v kafka +kafka=2.3.1 + +jmx java.lang:type=Memory +HeapMemoryUsage.init=1073741824;;NonHeapMemoryUsage.used=96265256;;NonHeapMemoryUsage.max=-1;;HeapMemoryUsage.max=1073741824;;NonHeapMemoryUsage.committed=101580800;;Verbose=true;;NonHeapMemoryUsage.init=7667712;;ObjectPendingFinalizationCount=0;;type=Memory;;HeapMemoryUsage.used=468633592;;HeapMemoryUsage.committed=1073741824 + +non-existing-command + +groups +group=mange1;;topic=pacman1;;partition=0;;current_offset=4003;;clientid=rdkafka;;consumerid=rdkafka-fb4bde72-24ff-4bd3-a2bc-d99267340e00;;host=/192.168.1.42 +group=mange2;;topic=pacman2;;partition=0;;current_offset=4003;;clientid=rdkafka;;consumerid=rdkafka-fb4bde72-24ff-4bd3-a2bc-d99267340e00;;host=/192.168.1.42 + +groups mange2 +group=mange2;;topic=pacman2;;partition=0;;current_offset=4003;;clientid=rdkafka;;consumerid=rdkafka-fb4bde72-24ff-4bd3-a2bc-d99267340e00;;host=/192.168.1.42 + +``` + +Since a command can give muliple results, it will always end with a empty row, +so for each command you give, make sure to read until empty row. + +If you give a command that doesn't exists or a JMX query that gives no results +it will just print a `\n` ## Usage -* Download the latest version from the releases -* Put the jar file in the `libs/` folder of your kafka deployment -* Add this line to `server.properties`: `metric.reporters=cloudkarafka.kafka_http_reporter` -* (Re)start the broker +- Download the latest version from the releases +- Put the jar file in the `libs/` folder of your kafka deployment +- Add this line to `server.properties`: `metric.reporters=cloudkarafka.kafka_http_reporter` +- (Re)start the broker ## Development -* `lein uberjar` produces a standalone jar file in `target/` -* Follow the same procedure as above to install it +- `lein uberjar` produces a standalone jar file in `target/` +- Follow the same procedure as above to install it ## Release -* Update the version number in `project.clj` -* Commit the change and tag it with the same version -* Push to github, travis CI will generate the release artifact +- Update the version number in `project.clj` +- Commit the change and tag it with the same version +- Push to github, travis CI will generate the release artifact ## Versioning -We use [SemVer](http://semver.org/) for versioning. For the versions available, see the [tags on this repository](https://github.com/CloudKarafka/KafkaHttpReporter/tags). +We use [SemVer](http://semver.org/) for versioning. For the versions available, see the [tags on this repository](https://github.com/CloudKarafka/KafkaHttpReporter/tags). ## Authors -* **Magnus Landerblom** - *Initial work* - [snichme](https://github.com/snichme) - - +- **Magnus Landerblom** - _Initial work_ - [snichme](https://github.com/snichme) diff --git a/src/cloudkarafka/cmds.clj b/src/cloudkarafka/cmds.clj index 479ff8d..72d8be9 100644 --- a/src/cloudkarafka/cmds.clj +++ b/src/cloudkarafka/cmds.clj @@ -1,6 +1,7 @@ (ns cloudkarafka.cmds (:require [clojure.java.jmx :as jmx] [cloudkarafka.kafkaadmin :as ka] + [cloudkarafka.util :as util] [clojure.string :as str]) (:import [javax.management ObjectName] [java.util Map$Entry])) @@ -59,15 +60,17 @@ (when bean (query bean)) (catch javax.management.OperationsException e - (println e)))) + (println "[WARN] KafkaHttpReporter jmx_error " (str e) (.getMessage e))))) (defmethod exec "v" [_ bean] (case bean - "kafka" (list {"kafka" (org.apache.kafka.common.utils.AppInfoParser/getVersion)}) - :else nil)) + "kafka" (list {"kafka" util/kafka-version}) + nil)) -(defmethod exec "groups" [_ group] - (ka/consumers group)) +(defmethod exec "groups" [_ groups] + (if (seq groups) + (ka/consumers (str/split groups #",")) + (ka/consumers))) (defmethod exec :default [_ _] nil) diff --git a/src/cloudkarafka/kafka_http_reporter.clj b/src/cloudkarafka/kafka_http_reporter.clj index 5c53077..d5488a8 100644 --- a/src/cloudkarafka/kafka_http_reporter.clj +++ b/src/cloudkarafka/kafka_http_reporter.clj @@ -24,7 +24,7 @@ {:status 200 :headers {"content-type" "text/plain"} :body "0.1.0"}) (GET "/kafka-version" [] - {:status 200 :headers {"content-type" "text/plain"} :body (:kafka-version @util/state)}) + {:status 200 :headers {"content-type" "text/plain"} :body (util/kafka-version)}) (GET "/jmx" [bean group attrs] (if-let [values (jmx/query bean group (str/split attrs #","))] @@ -58,20 +58,20 @@ [cmd bean] (str/split s #" ")] (when-not (empty? cmd) (let [res (exec cmd bean)] - (str (when res (format-result res)) "\n"))))) + (str (when (seq res) (format-result res)) "\n"))))) + +(defn props-from-config [config] + (let [ listener-name (or (:security.inter.broker.protocol config) "PLAINTEXT") + uris (util/listener-uri (:listeners config) listener-name)] + {:bootstrap.servers (first uris)})) (defn -configure [_ config] (println "[INFO] KafkaHttpReporter: configure") (let [parsed-config (into {} (map (fn [[k v]] [(keyword k) v]) config)) - listener-name (or (:security.inter.broker.protocol parsed-config) "PLAINTEXT") - uris (util/listener-uri (:listeners parsed-config) listener-name) - props {:bootstrap.servers (first uris)} - kafka-version (org.apache.kafka.common.utils.AppInfoParser/getVersion)] - (reset! util/state {:kafka-version kafka-version - :kafka-config parsed-config - :admin-client (when (util/modern-kafka? kafka-version) - (ka/admin-client props)) - :consumer (ka/kafka-consumer props)}))) + props (props-from-config parsed-config)] + (reset! util/state {:kafka-config parsed-config + :admin-client (when (util/modern-kafka?) + (ka/admin-client props))}))) (defn -init [_ _] (let [config (:kafka-config @util/state) diff --git a/src/cloudkarafka/kafkaadmin.clj b/src/cloudkarafka/kafkaadmin.clj index e5c5154..3885107 100644 --- a/src/cloudkarafka/kafkaadmin.clj +++ b/src/cloudkarafka/kafkaadmin.clj @@ -10,45 +10,36 @@ (.put props (name k) v)) props)) -(defn member-list [group-name desc group-offset log-offset] +(defn member-list [group-name desc group-offset] (for [member (.members desc) :let [ toppar (.topicPartitions (.assignment member)) ]] (if (empty? toppar) {:group group-name - :topic nil, - :partition nil, - :current_offset nil, - :log_end_offset nil, - :lag nil, :clientid (.clientId member), :consumerid (.consumerId member), :host (.host member)} (for [tp toppar :let [group-partitions (get group-offset tp) - go (when group-partitions (.offset group-partitions)) - lo (get log-offset tp)]] + go (when group-partitions (.offset group-partitions))]] {:group group-name :topic (.topic tp) :partition (.partition tp) :current_offset go, - :log_end_offset lo, - :lag (when go (- lo go)), :clientid (.clientId member), :consumerid (.consumerId member), :host (.host member)})))) (defn consumer-groups - ([client consumer] - (consumer-groups client consumer nil)) - ([client consumer group-ids] + ([client] + (consumer-groups client nil)) + ([client group-ids] (let [group-ids (or group-ids (into [] (map #(.groupId %) (.get (.all (.listConsumerGroups client)))))) descs (.get (.all (.describeConsumerGroups client group-ids)))] (flatten (for [group-id group-ids :let [desc (get descs group-id) - group-offset (.get (.partitionsToOffsetAndMetadata (.listConsumerGroupOffsets client group-id))) - log-offset (.endOffsets consumer (mapv key group-offset))]] - (member-list group-id desc group-offset log-offset)))))) + group-offset (.get (.partitionsToOffsetAndMetadata (.listConsumerGroupOffsets client group-id)))]] + (member-list group-id desc group-offset)))))) (defn filtered-groups [client group-ids] (let [all-group-ids (map #(.groupId %) (.listAllConsumerGroupsFlattened client)) @@ -58,25 +49,21 @@ (filter #(contains? wanted %) all-group-ids)))) (defn consumer-groups-old - ([url consumer group-ids] + ([url group-ids] (with-open [client (kafka.admin.AdminClient/createSimplePlaintext url)] (let [res (java.util.LinkedList.)] ($/for [group-id (filtered-groups client group-ids) :let [summary (.describeConsumerGroup client group-id 0) group-offset (.listGroupOffsets client group-id)]] ($/if-let [members (.consumers summary)] - ($/for [member members - :let [log-end-offsets (.endOffsets consumer (scala.collection.JavaConversions/asJavaCollection (.assignment member)))]] + ($/for [member members] ($/for [toppar (.assignment member) - :let [current-offset (get (scala.collection.JavaConversions/mapAsJavaMap group-offset) toppar) - log-end (get log-end-offsets toppar)]] + :let [current-offset (get (scala.collection.JavaConversions/mapAsJavaMap group-offset) toppar)]] (.add res {:state (.state summary) :group group-id :topic (.topic toppar) :partition (.partition toppar) :current_offset current-offset - :log_end_offset log-end - :lag (and log-end current-offset (- log-end current-offset)) :consumerid (.consumerId member) :clientid (.clientId member) :host (.host member)}))) @@ -87,35 +74,17 @@ [props] (org.apache.kafka.clients.admin.AdminClient/create (map->props props))) -(defn kafka-consumer ^org.apache.kafka.clients.consumer.KafkaConsumer - [props] - (let [m (merge props - {:group.id "mgmt-admin", - :key.deserializer "org.apache.kafka.common.serialization.StringDeserializer", - :value.deserializer "org.apache.kafka.common.serialization.StringDeserializer"}) - props (map->props m)] - (org.apache.kafka.clients.consumer.KafkaConsumer. props))) - (defn consumers ([] (consumers nil)) - ([group] + ([groups] (let [s @util/state kafka-config (:kafka-config s)] - (if (util/modern-kafka? (:kafka-version s)) - (consumer-groups (:admin-client s) (:consumer s) group) + (if (util/modern-kafka?) + (consumer-groups (:admin-client s) groups) (let [plaintext-url (-> kafka-config :listeners (util/listener-uri (or (:security.inter.broker.protocol kafka-config) "PLAINTEXT")) first)] - (consumer-groups-old plaintext-url (:consumer s) group)))))) - -(comment - (consumer-groups-old - "127.0.0.1:9092" - (kafka-consumer {:bootstrap.servers "127.0.0.1:9092"})) + (consumer-groups-old plaintext-url (:consumer s) groups)))))) - (consumer-groups - (admin-client {:bootstrap.servers "127.0.0.1:9092"}) - (kafka-consumer {:bootstrap.servers "127.0.0.1:9092"})) -) diff --git a/src/cloudkarafka/util.clj b/src/cloudkarafka/util.clj index 81517d5..2ba78b5 100644 --- a/src/cloudkarafka/util.clj +++ b/src/cloudkarafka/util.clj @@ -6,15 +6,19 @@ (def state (atom nil)) +(def kafka-version ( org.apache.kafka.common.utils.AppInfoParser/getVersion)) + (defn listener-uri [listeners type] (->> (str/split listeners #",") (map #(str/split % #"://")) (filter #(= (first %) type)) (map second))) -(defn modern-kafka? [version] - (let [major (-> version - (str/split #"\.") - first - Integer/parseInt)] - (>= major 2))) +(defn modern-kafka? + ([] (modern-kafka? kafka-version)) + ([version] + (let [major (-> version + (str/split #"\.") + first + Integer/parseInt)] + (>= major 2)))) From 52985a5ce60f8b014e9c9ac9a8f62d9900dc9e59 Mon Sep 17 00:00:00 2001 From: Magnus Landerblom Date: Fri, 17 Apr 2020 20:38:59 +0200 Subject: [PATCH 6/6] bug --- src/cloudkarafka/kafka_http_reporter.clj | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/src/cloudkarafka/kafka_http_reporter.clj b/src/cloudkarafka/kafka_http_reporter.clj index d5488a8..a662d5a 100644 --- a/src/cloudkarafka/kafka_http_reporter.clj +++ b/src/cloudkarafka/kafka_http_reporter.clj @@ -24,7 +24,7 @@ {:status 200 :headers {"content-type" "text/plain"} :body "0.1.0"}) (GET "/kafka-version" [] - {:status 200 :headers {"content-type" "text/plain"} :body (util/kafka-version)}) + {:status 200 :headers {"content-type" "text/plain"} :body util/kafka-version}) (GET "/jmx" [bean group attrs] (if-let [values (jmx/query bean group (str/split attrs #","))]