コンテンツにスキップ

Kafka 4

Kafka 4 のメトリクスデータを収集

設定

前提条件

JMX Exporter をダウンロード

ダウンロード先:https://github.com/prometheus/jmx_exporter/releases/tag/1.3.0

JMX スクリプトと起動引数を設定

注意:Producer、Consumer、Streams、Connect のメトリクスを収集するには、それぞれ独立したプロセスを起動する必要があります。各プロセスを起動する際は、対応する yaml ファイルと起動スクリプトを置き換えてください。以下を参考にできます。

KRaft Metrics
  • KRaft Metrics の設定ファイル kafka.yml を作成
# ------------------------------------------------------------
# Kafka 4 Prometheus JMX Exporter Configuration
# ------------------------------------------------------------
lowercaseOutputName: false
lowercaseOutputLabelNames: true
cacheRules: true
rules:

# 1. Broker / Topic / Partition Metrics
  - pattern: kafka.server<type=BrokerTopicMetrics, name=(BytesInPerSec|BytesOutPerSec|MessagesInPerSec|TotalFetchRequestsPerSec|ProduceRequestsPerSec|FailedProduceRequestsPerSec|TotalProduceRequestsPerSec|ReassignmentBytesInPerSec|ReassignmentBytesOutPerSec|ProduceMessageConversionsPerSec|FetchMessageConversionsPerSec)(?:, topic=([-\.\w]*))?><>(Count|OneMinuteRate|FiveMinuteRate|FifteenMinuteRate|MeanRate)
    name: kafka_server_broker_topic_metrics_$1
    type: GAUGE
    labels:
      topic: "$2"

# 2. Request / Network Metrics
  - pattern: kafka.network<type=RequestMetrics, name=(.+)><>(Count|OneMinuteRate|FiveMinuteRate|FifteenMinuteRate|MeanRate)
    name: kafka_network_request_metrics_$1
    type: GAUGE

# 3. Socket Server Metrics
  - pattern: kafka.network<type=SocketServer, name=(.+)><>(Count|OneMinuteRate|FiveMinuteRate|FifteenMinuteRate|MeanRate|Value)
    name: kafka_network_socket_server_metrics_$1
    type: GAUGE

# 4. Log / Segment / Cleaner Metrics
  - pattern: kafka.log<type=LogFlushStats, name=(.+)><>(Count|OneMinuteRate|FiveMinuteRate|FifteenMinuteRate|MeanRate)
    name: kafka_log_$1_$2
    type: GAUGE

# 5. Controller (KRaft) Metrics
  - pattern: kafka.controller<type=KafkaController, name=(.+)><>(Count|Value)
    name: kafka_controller_$1
    type: GAUGE

# 6. Group / Coordinator Metrics
  - pattern: kafka.coordinator.group<type=GroupMetadataManager, name=(.+)><>(Count|Value)
    name: kafka_coordinator_group_metadata_manager_$1
    type: GAUGE

# 7. KRaft Specific Metrics
  - pattern: kafka.controller<type=KafkaController, name=(LeaderElectionSuccessRate|LeaderElectionLatencyMs)><>(Count|Value)
    name: kafka_controller_$1
    type: GAUGE

# 8. New Generation Consumer Rebalance Protocol Metrics
  - pattern: kafka.coordinator.group<type=GroupMetadataManager, name=(RebalanceTimeMs|RebalanceFrequency)><>(Count|Value)
    name: kafka_coordinator_group_metadata_manager_$1
    type: GAUGE

# 9. Queue Metrics
  - pattern: kafka.server<type=Queue, name=(QueueSize|QueueConsumerRate)><>(Count|Value)
    name: kafka_server_queue_$1
    type: GAUGE

# 10. Client Metrics
  - pattern: kafka.network<type=RequestMetrics, name=(ClientConnections|ClientRequestRate|ClientResponseTime)><>(Count|OneMinuteRate|FiveMinuteRate|FifteenMinuteRate|MeanRate)
    name: kafka_network_request_metrics_$1
    type: GAUGE

# 11. Log Flush Rate and Time
  - pattern: kafka.log<type=LogFlushStats, name=LogFlushRateAndTimeMs><>(Count|OneMinuteRate|FiveMinuteRate|FifteenMinuteRate|MeanRate)
    name: kafka_log_log_flush_rate_and_time_ms
    type: GAUGE
  • 起動引数
export KAFKA_HEAP_OPTS="-Xms1g -Xmx1g"
export KAFKA_JMX_OPTS="-Dcom.sun.management.jmxremote.port=9999 \
  -Dcom.sun.management.jmxremote.rmi.port=9999 \
  -Dcom.sun.management.jmxremote.authenticate=false \
  -Dcom.sun.management.jmxremote.ssl=false \
  -Djava.rmi.server.hostname=127.0.0.1"
export Environment="KAFKA_OPTS=-javaagent:/opt/jmx_exporter/jmx_prometheus_javaagent.jar=7071:/opt/jmx_exporter/kafka.yml"

/opt/kafka/bin/kafka-server-start.sh -daemon /opt/kafka/config/kraft/server.properties
Producer Metrics
  • Producer Metrics の設定ファイル producer.yml を作成
---
lowercaseOutputName: true
rules:
  # 新規追加:producer-node-metrics
  - pattern: kafka\.producer<type=producer-node-metrics, client-id=([^,]+), node-id=([^>]+)><>([^:]+)
    name: kafka_producer_node_$3
    labels:
      client_id: "$1"
      node_id: "$2"
    type: GAUGE

  - pattern: 'kafka\.producer<type=producer-metrics, client-id=([^>]+)><>([^:,\s]+).*'
    name: 'kafka_producer_metrics_$2'
    labels:
      client_id: "$1"
    type: GAUGE

  # Selector の全メトリクスを取得(Kafka 4.0 で追加)
  - pattern: 'kafka\.(?:(producer|consumer|connect))<type=(producer|consumer|connect)-metrics, client-id=([^>]+)><>(connection-.+|io-.+|network-.+|select-.+|send-.+|receive-.+|reauthentication-.+)'
    name: 'kafka_${1}_${4}'
    labels:
      client_id: '$3'
    type: GAUGE
  • 起動引数
export Environment="KAFKA_OPTS=-javaagent:/opt/jmx_exporter/jmx_prometheus_javaagent.jar=7072:/opt/jmx_exporter/producer.yml"

/opt/kafka/kafka/bin/kafka-console-producer.sh \
  --broker-list localhost:9092 \
  --topic xxxx \
  --producer-property bootstrap.servers=localhost:9092
Consumer Metrics
  • Consumer Metrics の設定ファイル consumer.yml を作成
lowercaseOutputName: true
rules:
  # consumer-coordinator-metrics
  - pattern: 'kafka\.consumer<type=consumer-coordinator-metrics, client-id=([^>]+)><>([^:,\s]+).*'
    name: 'kafka_consumer_coordinator_metrics_$2'
    labels:
      client_id: "$1"
    type: GAUGE

  - pattern: 'kafka\.consumer<type=consumer-metrics, client-id=([^>]+)><>([^:,\s]+).*'
    name: 'kafka_consumer_metrics_$2'
    labels:
      client_id: "$1"
  • 起動引数
export Environment="KAFKA_OPTS=-javaagent:/opt/jmx_exporter/jmx_prometheus_javaagent.jar=7073:/opt/jmx_exporter/consumer.yml"

/opt/kafka/kafka/bin/kafka-console-consumer.sh \
  --broker-list localhost:9092 \
  --topic xxxx \
  --producer-property bootstrap.servers=localhost:9092
Streams Metrics
  • Streams Metrics の設定ファイル stream.yml を作成
lowercaseOutputName: true
lowercaseOutputLabelNames: true

rules:
  # Kafka Streams アプリケーション指標 - 特殊文字を削除
  - pattern: 'kafka.streams<type=stream-metrics, client-id=(.+)><>([a-zA-Z0-9\-]+)$'
    name: kafka_streams_$2
    labels:
      client_id: "$1"

  # 特殊文字を含む属性名を処理
  - pattern: 'kafka.streams<type=stream-metrics, client-id=(.+)><>([a-zA-Z0-9\-]+):(.+)$'
    name: kafka_streams_$2_$3
    labels:
      client_id: "$1"

  # Processor Node 指標
  - pattern: 'kafka.streams<type=stream-processor-node-metrics, client-id=(.+), task-id=(.+), processor-node-id=(.+)><>(.+)'
    name: kafka_streams_processor_$4
    labels:
      client_id: "$1"
      task_id: "$2"
      processor_node_id: "$3"

  # Task 指標
  - pattern: 'kafka.streams<type=stream-task-metrics, client-id=(.+), task-id=(.+)><>(.+)'
    name: kafka_streams_task_$3
    labels:
      client_id: "$1"
      task_id: "$2"

  # スレッド指標
  - pattern: 'kafka.streams<type=stream-thread-metrics, client-id=(.+), thread-id=(.+)><>(.+)'
    name: kafka_streams_thread_$3
    labels:
      client_id: "$1"
      thread_id: "$2"

  # JVM 指標
  - pattern: 'java.lang<type=Memory><>(.+)'
    name: jvm_memory_$1

  - pattern: 'java.lang<type=GarbageCollector, name=(.+)><>(\w+)'
    name: jvm_gc_$2
    labels:
      gc: "$1"

  # スレッドプール指標
  - pattern: 'java.lang<type=Threading><>(.+)'
    name: jvm_threads_$1

  # デフォルトルール
  - pattern: '(.*)'
  • 起動引数
export KAFKA_HEAP_OPTS="-Xms512m -Xmx512m"
export KAFKA_JMX_OPTS="-Dcom.sun.management.jmxremote.port=9996 \
  -Dcom.sun.management.jmxremote.rmi.port=9996 \
  -Dcom.sun.management.jmxremote.authenticate=false \
  -Dcom.sun.management.jmxremote.ssl=false \
  -Djava.rmi.server.hostname=127.0.0.1"
export Environment="KAFKA_OPTS=-javaagent:/opt/jmx_exporter/jmx_prometheus_javaagent.jar=7075:/opt/jmx_exporter/stream.yml"

java $KAFKA_HEAP_OPTS $KAFKA_JMX_OPTS $EXTRA_ARGS -cp "libs/*:my-streams.jar" WordCountDemo
Connect Metrics
  • Connect Metrics の設定ファイル connect.yml を作成
lowercaseOutputName: true
lowercaseOutputLabelNames: true
rules:
  # 1) connect-worker-metrics(グローバル)
  - pattern: 'kafka\.connect<type=connect-worker-metrics><>([^:]+)'
    name: 'kafka_connect_worker_$1'
    type: GAUGE

  # 2) connect-worker-metrics,connector=xxx
  - pattern: 'kafka\.connect<type=connect-worker-metrics, connector=([^>]+)><>([^:]+)'
    name: 'kafka_connect_worker_$2'
    labels:
      connector: "$1"
    type: GAUGE

  # 3) connect-worker-rebalance-metrics
  - pattern: 'kafka\.connect<type=connect-worker-rebalance-metrics><>([^:]+)'
    name: 'kafka_connect_worker_rebalance_$1'
    type: GAUGE

  # 4) connector-task-metrics
  - pattern: 'kafka\.connect<type=connector-task-metrics, connector=([^>]+), task=([^>]+)><>([^:]+)'
    name: 'kafka_connect_task_$3'
    labels:
      connector: "$1"
      task_id: "$2"
    type: GAUGE

  # 5) sink-task-metrics
  - pattern: 'kafka\.connect<type=sink-task-metrics, connector=([^>]+), task=([^>]+)><>([^:]+)'
    name: 'kafka_connect_sink_task_$3'
    labels:
      connector: "$1"
      task_id: "$2"
  • 起動引数
export KAFKA_HEAP_OPTS="-Xms512m -Xmx512m"
export KAFKA_JMX_OPTS="-Dcom.sun.management.jmxremote \
  -Dcom.sun.management.jmxremote.authenticate=false \
  -Dcom.sun.management.jmxremote.ssl=false \
  -Dcom.sun.management.jmxremote.port=9995 \
  -Dcom.sun.management.jmxremote.rmi.port=9995 \
  -Djava.rmi.server.hostname=127.0.0.1"
export Environment="KAFKA_OPTS=-javaagent:/opt/jmx_exporter/jmx_prometheus_javaagent.jar=7074:/opt/jmx_exporter/connect.yml"

# Kafka Connect を起動
/opt/kafka/kafka/bin/connect-distributed.sh /opt/kafka/kafka/config/connect-distributed.properties

起動成功後、curl http://IP:ポート番号/metrics で取得した監視データを確認できます。

DataKit を設定

  • datakit のインストールディレクトリ内の conf.d/samples ディレクトリへ移動し、prom.conf.sample をコピーして kafka.conf という名前に変更

cp prom.conf.sample kafka.conf

  • kafka.conf を調整
[[inputs.prom]]
  ## Exporter URLs.
  urls = ["http://127.0.0.1:7071/metrics","http://127.0.0.1:7072/metrics","http://127.0.0.1:7073/metrics","http://127.0.0.1:7074/metrics","http://127.0.0.1:7075/metrics"]

  ## Collector alias.
  source = "kafka"

  ## Prioritier over 'measurement_name' configuration.
  [[inputs.prom.measurements]]
    prefix = "kafka_controller_"
    name = "kafka_controller"

  [[inputs.prom.measurements]]
    prefix = "kafka_network_"
    name = "kafka_network"

  [[inputs.prom.measurements]]
    prefix = "kafka_log_"
    name = "kafka_log"

  [[inputs.prom.measurements]]
    prefix = "kafka_server_"
    name = "kafka_server"

  [[inputs.prom.measurements]]
    prefix = "kafka_connect_"
    name = "kafka_connect"

  [[inputs.prom.measurements]]
    prefix = "kafka_stream_"
    name = "kafka_stream"
  • DataKit を再起動

以下のコマンドを実行

datakit service -R

メトリクス

以下は kafka4 の一部メトリクスの説明です。より多くのメトリクスは Kafka 指標詳細 を参照してください

kafka_server メトリクスセット

指標名 説明 単位
Fetch_queue_size Fetch キューサイズ count
Produce_queue_size Producer キューサイズ count
Request_queue_size Request キューサイズ count
broker_topic_metrics_BytesInPerSec クライアントの受信バイトレート bytes/s
broker_topic_metrics_BytesOutPerSec クライアントの送信バイトレート bytes/s
broker_topic_metrics_FailedProduceRequestsPerSec 送信リクエスト失敗率 count/s
broker_topic_metrics_FetchMessageConversionsPerSec Fetch メッセージ変換レート count/s
broker_topic_metrics_MessagesInPerSec 受信メッセージレート count/s
broker_topic_metrics_ProduceMessageConversionsPerSec Producer メッセージ変換レート count/s
broker_topic_metrics_TotalFetchRequestsPerSec フェッチリクエスト総レート(クライアントまたは follower から) count/s
broker_topic_metrics_TotalProduceRequestsPerSec Producer リクエストレート count/s
socket_server_metrics_connection_count SocketServer の接続数 count
socket_server_metrics_connection_close_total SocketServer の切断済み接続数 count
socket_server_metrics_incoming_byte_rate SocketServer の受信バイトレート bytes/s

kafka_network メトリクスセット

指標名 説明 単位
request_metrics_RequestBytes_request_AddOffsetsToTxn AddOffsetsToTxn リクエストサイズ bytes
request_metrics_RequestBytes_request_Fetch Fetch リクエストサイズ count
request_metrics_RequestBytes_request_FetchConsumer FetchConsumer リクエストサイズ bytes
request_metrics_RequestBytes_request_FetchFollower FetchFollower リクエストサイズ bytes
request_metrics_TotalTimeMs_request_CreateTopics CreateTopics リクエストの総時間 ms
request_metrics_TotalTimeMs_request_CreatePartitions CreatePartitions リクエストの総時間 ms
request_metrics_RequestQueueTimeMs_request_CreateTopics CreateTopics のリクエストキュー待ち時間 ms
request_metrics_RequestQueueTimeMs_request_CreatePartitions CreatePartitions のリクエストキュー待ち時間 ms
request_metrics_RequestQueueTimeMs_request_Produce Produce のリクエストキュー待ち時間 ms
request_metrics_ResponseSendTimeMs_request_CreateTopics CreateTopics のレスポンス送信時間 ms
request_metrics_ResponseSendTimeMs_request_CreatePartitions CreatePartitions のレスポンス送信時間 ms

kafka_controller メトリクスセット

指標名 説明 単位
ActiveBrokerCount アクティブな Broker 数 count
ActiveControllerCount アクティブなコントローラ数 count
GlobalPartitionCount パーティション数 count
GlobalTopicCount トピック数 count
OfflinePartitionsCount オフラインパーティション数 count
PreferredReplicaImbalanceCount Preferred Leader 選出条件を満たすパーティション数 count
OfflinePartitionsCount オフラインパーティション数 count
TimedOutBrokerHeartbeatCount Broker ハートビートのタイムアウト回数 count
LastAppliedRecordLagMs 最後に適用されたレコードの遅延時間 ms
LastAppliedRecordOffset 最後に適用されたレコードのオフセット -
MetadataErrorCount メタデータエラー数 count
NewActiveControllersCount 新しいコントローラ選出回数 count

kafka_producer メトリクスセット

指標名 説明 単位
producer_metrics_batch_split_rate バッチ分割率 count/s
producer_metrics_buffer_available_bytes 未使用バッファメモリ総量 bytes
producer_metrics_buffer_exhausted_rate バッファ枯渇によって破棄された 1 秒あたりの平均送信レコード数 count/s
producer_metrics_buffer_total_bytes バッファ総バイト数 bytes
producer_metrics_bufferpool_wait_ratio バッファプール待機比率 %
producer_metrics_bufferpool_wait_time_ns_total バッファプール待機時間 ms
producer_metrics_connection_close_rate 接続クローズ率 count/s
producer_metrics_connection_count クローズされた接続数 count
producer_metrics_flush_time_ns_total フラッシュ総時間 ns
producer_metrics_incoming_byte_rate 受信バイト率 bytes/s
producer_metrics_outgoing_byte_rate 送信バイト率 bytes/s
producer_metrics_request_rate リクエスト率 count/s
producer_metrics_request_size_avg リクエストサイズ平均 bytes

kafka_consumer メトリクスセット

指標名 説明 単位
consumer_coordinator_metrics_failed_rebalance_total 再バランス失敗数 count
consumer_coordinator_metrics_heartbeat_rate 1 秒あたりの平均ハートビート回数 count/s
consumer_coordinator_metrics_heartbeat_response_time_max ハートビート応答最大時間 count
consumer_coordinator_metrics_join_rate Group 参加レート count/s
consumer_coordinator_metrics_join_total Group 参加総数 count
consumer_coordinator_metrics_last_rebalance_seconds_ago 前回の再バランスイベントからの経過秒数 ms
consumer_coordinator_metrics_rebalance_latency_total 再バランス遅延合計 ms
consumer_fetch_manager_metrics_bytes_consumed_rate 1 秒あたりの消費バイト数 bytes/s
consumer_fetch_manager_metrics_fetch_latency_avg Fetch リクエスト遅延 ms
consumer_metrics_connection_count 接続数 count
consumer_metrics_connection_count 切断された接続数 count/s
consumer_metrics_incoming_byte_rate 受信バイト率 bytes/s
consumer_metrics_outgoing_byte_rate 送信バイト率 bytes/s
consumer_metrics_select_rate Select レート count/s
consumer_metrics_last_poll_seconds_ago I/O 待機時間 ms
consumer_metrics_last_poll_seconds_ago I/O 待機時間 ms

kafka_connect メトリクスセット

指標名 説明 単位
worker_connector_count Connector 数 count
worker_task_startup_attempts_total タスク起動再試行回数 count
worker_connector_startup_attempts_total Connector 起動試行回数 count
worker_task_startup_failure_total タスク起動失敗数 count
worker_connector_startup_failure_percentage Connector 起動失敗率 %
worker_rebalance_completed_rebalances_total 再バランス完了総数 count
worker_task_startup_failure_percentage タスク起動失敗率 %
worker_rebalance_time_since_last_rebalance_ms 前回の再バランスからの経過時間 ms
worker_task_startup_attempts_total タスク起動試行回数 count

kafka_stream メトリクスセット

指標名 説明 単位
stream_thread_metrics_thread_start_time スレッド起動時間 時間戳 ms
stream_thread_metrics_task_created_total タスク作成総数 count
stream_state_metrics_block_cache_capacity ブロックキャッシュサイズ bytes
stream_state_metrics_all_rate 全操作レート count/s
stream_state_metrics_block_cache_usage ブロックキャッシュ使用率 %
stream_state_metrics_bytes_read_compaction_rate バイト読み取り圧縮率 bytes/s
stream_state_metrics_bytes_written_compaction_rate バイト書き込み圧縮率 bytes/s
stream_state_metrics_block_cache_index_hit_ratio ブロックキャッシュ索引ヒット率 %
stream_state_metrics_block_cache_data_hit_ratio ブロックキャッシュデータヒット率 %
stream_state_metrics_block_cache_filter_hit_ratio ブロックキャッシュフィルターヒット率 %
stream_state_metrics_bytes_written_rate バイト書き込みレート bytes/s
stream_state_metrics_bytes_read_rate バイト読み取りレート bytes/s
stream_state_metrics_block_cache_filter_hit_ratio キャッシュサイズのバイト数 bytes
stream_task_metrics_process_rate 1 秒あたりの処理レコード数 bytes/s
stream_task_metrics_enforced_processing_rate 1 秒あたりの強制処理数 bytes/s
stream_task_metrics_active_process_ratio アクティブ処理比率 %
stream_thread_metrics_commit_rate 提出率 count/s
stream_thread_metrics_poll_latency_avg ポーリング遅延時間 ms
stream_thread_metrics_poll_rate ポーリングレート count/s
stream_thread_metrics_blocked_time_ns_total ブロック時間 ns
stream_topic_metrics_bytes_consumed_total 消費バイト数 bytes
stream_topic_metrics_bytes_produced_total 生成バイト数 bytes
stream_topic_metrics_records_consumed_total ソースプロセッサノードが消費したレコード総数 count
stream_topic_metrics_records_produced_total シンクプロセッサノードが生成したレコード総数 count

フィードバック

このページは役に立ちましたか?