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/project.clj b/project.clj index 7b91a95..e102b67 100644 --- a/project.clj +++ b/project.clj @@ -1,14 +1,15 @@ -(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} :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/cmds.clj b/src/cloudkarafka/cmds.clj new file mode 100644 index 0000000..72d8be9 --- /dev/null +++ b/src/cloudkarafka/cmds.clj @@ -0,0 +1,78 @@ +(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])) + +(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 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 mbean-params [^ObjectName mbean] + (let [kpl (.getKeyPropertyList mbean) + entries (.entrySet kpl)] + (into {} (map mbean-keys->map entries)))) + +(defn jmx-values [mbean] + (let [params (mbean-params mbean) + values (build-value-map "" (jmx/mbean mbean))] + (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 format-result-row [values] + (->> values + (map (fn [[k v]] (str (name k) "=" v))) + (str/join ";;"))) + +(defn format-result [values] + (str + (str/join "\n" (map format-result-row values)) + "\n")) + +(defmulti exec (fn [a _] a)) + +(defmethod exec "jmx" [_ bean] + (try + (when bean + (query bean)) + (catch javax.management.OperationsException e + (println "[WARN] KafkaHttpReporter jmx_error " (str e) (.getMessage e))))) + +(defmethod exec "v" [_ bean] + (case bean + "kafka" (list {"kafka" util/kafka-version}) + nil)) + +(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 9d1c7a6..a662d5a 100644 --- a/src/cloudkarafka/kafka_http_reporter.clj +++ b/src/cloudkarafka/kafka_http_reporter.clj @@ -1,8 +1,12 @@ (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] @@ -13,36 +17,6 @@ (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 @@ -50,52 +24,69 @@ {: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 util/kafka-version}) (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 -configure [this config] - (let [parsed-config (into {} (map (fn [[k v]] [(keyword k) v]) config)) - uris (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 - :kafka-config parsed-config - :admin-client (when (modern-kafka? kafka-version) - (ka/admin-client props)) - :consumer (ka/kafka-consumer props)}))) +(defn wrapper [f] + (fn [s _info] + (s/connect (s/map f s) s))) -(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})))) +(defn tcp-handler [^bytes in] + (let [s (str/trim (String. in)) + [cmd bean] (str/split s #" ")] + (when-not (empty? cmd) + (let [res (exec cmd bean)] + (str (when (seq res) (format-result res)) "\n"))))) -(defn -metricChange [this metric]) -(defn -metricRemoval [this metric]) +(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 -close [this] - (when-let [^java.io.Closeable s (:http-server @state)] - (println "[INFO] KafkaHttpReporter: Closing HTTP server") - (.close s))) +(defn -configure [_ config] + (println "[INFO] KafkaHttpReporter: configure") + (let [parsed-config (into {} (map (fn [[k v]] [(keyword k) v]) config)) + 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) + 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/start-server (wrapper tcp-handler) {:port tcp-port})))) + +(defn -metricChange [_ _]) +(defn -metricRemoval [_ _]) +(defn -close [_]) diff --git a/src/cloudkarafka/kafkaadmin.clj b/src/cloudkarafka/kafkaadmin.clj index f02d13e..3885107 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] @@ -9,84 +10,81 @@ (.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)]] - {: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))]] + {:group group-name + :topic (.topic tp) + :partition (.partition tp) + :current_offset 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-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)))]] + (member-list group-id desc group-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] - (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 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] + ($/for [toppar (.assignment member) + :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 + :consumerid (.consumerId member) + :clientid (.clientId member) + :host (.host member)}))) + :none)) + res)))) (defn admin-client ^org.apache.kafka.clients.admin.AdminClient [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))) - -(comment - - (consumer-groups-old "127.0.0.1:9092" (kafka-consumer {:bootstrap.servers "127.0.0.1:9092"})) +(defn consumers + ([] (consumers nil)) + ([groups] + (let [s @util/state + kafka-config (:kafka-config s)] + (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) groups)))))) - ) diff --git a/src/cloudkarafka/tcp.clj b/src/cloudkarafka/tcp.clj new file mode 100644 index 0000000..7dc054b --- /dev/null +++ b/src/cloudkarafka/tcp.clj @@ -0,0 +1,73 @@ +(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 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))) + +(defn valid-value? [v] + (or + (string? v) + (int? v) + (double? v) + (boolean? v))) + +(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 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] + (let [params (mbean-params mbean) + values (build-value-map "" (jmx/mbean mbean))] + (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))))) + +(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 {})) + + diff --git a/src/cloudkarafka/util.clj b/src/cloudkarafka/util.clj new file mode 100644 index 0000000..2ba78b5 --- /dev/null +++ b/src/cloudkarafka/util.clj @@ -0,0 +1,24 @@ +(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)) + +(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? + ([] (modern-kafka? kafka-version)) + ([version] + (let [major (-> version + (str/split #"\.") + first + Integer/parseInt)] + (>= major 2))))