From 01c2acd0b77b2819023c4825a2b39e3c8cc6c7b2 Mon Sep 17 00:00:00 2001 From: 3AceShowHand Date: Mon, 17 Aug 2026 15:28:55 +0800 Subject: [PATCH 1/2] add throttle time metrics and grafana --- metrics/grafana/ticdc_new_arch.json | 106 +++++++++++++++++- .../ticdc_new_arch_next_gen.json | 106 +++++++++++++++++- .../ticdc_new_arch_with_keyspace_name.json | 106 +++++++++++++++++- pkg/sink/kafka/metrics.go | 8 ++ pkg/sink/kafka/metrics_collector.go | 17 +++ 5 files changed, 340 insertions(+), 3 deletions(-) diff --git a/metrics/grafana/ticdc_new_arch.json b/metrics/grafana/ticdc_new_arch.json index 0a3f4bbd37..8dda9ef46a 100644 --- a/metrics/grafana/ticdc_new_arch.json +++ b/metrics/grafana/ticdc_new_arch.json @@ -19166,6 +19166,110 @@ "align": false, "alignLevel": null } + }, + { + "aliasColors": {}, + "bars": false, + "dashLength": 10, + "dashes": false, + "datasource": "${DS_TEST-CLUSTER}", + "description": "Average and p99 throttle time reported by each broker.", + "fieldConfig": { + "defaults": { + "links": [] + }, + "overrides": [] + }, + "fill": 1, + "fillGradient": 0, + "gridPos": { + "h": 7, + "w": 12, + "x": 12, + "y": 67 + }, + "hiddenSeries": false, + "id": 62108, + "legend": { + "alignAsTable": true, + "avg": false, + "current": true, + "max": true, + "min": false, + "rightSide": false, + "show": true, + "total": false, + "values": true + }, + "lines": true, + "linewidth": 1, + "links": [], + "nullPointMode": "null", + "options": { + "alertThreshold": true + }, + "paceLength": 10, + "percentage": false, + "pluginVersion": "7.5.17", + "pointradius": 2, + "points": false, + "renderer": "flot", + "seriesOverrides": [], + "spaceLength": 10, + "stack": false, + "steppedLine": false, + "targets": [ + { + "exemplar": true, + "expr": "sum(ticdc_sink_kafka_producer_throttle_time{k8s_cluster=\"$k8s_cluster\", tidb_cluster=\"$tidb_cluster\", namespace=~\"$namespace\",changefeed=~\"$changefeed\", instance=~\"$ticdc_instance\"}) by (namespace,changefeed, instance, broker, type)", + "format": "time_series", + "interval": "", + "intervalFactor": 1, + "legendFormat": "{{namespace}}-{{changefeed}}-{{instance}}-{{broker}}-{{type}}", + "refId": "A" + } + ], + "thresholds": [], + "timeFrom": null, + "timeRegions": [], + "timeShift": null, + "title": "Kafka Throttle Time", + "tooltip": { + "shared": true, + "sort": 0, + "value_type": "individual" + }, + "type": "graph", + "xaxis": { + "buckets": null, + "mode": "time", + "name": null, + "show": true, + "values": [] + }, + "yaxes": [ + { + "decimals": 1, + "format": "ms", + "label": null, + "logBase": 1, + "max": null, + "min": "0", + "show": true + }, + { + "format": "short", + "label": null, + "logBase": 1, + "max": null, + "min": null, + "show": false + } + ], + "yaxis": { + "align": false, + "alignLevel": null + } } ], "title": "Sink - MQ Sink", @@ -28342,5 +28446,5 @@ "timezone": "browser", "title": "${DS_TEST-CLUSTER}-TiCDC-New-Arch", "uid": "YiGL8hBZ0aac", - "version": 41 + "version": 42 } diff --git a/metrics/nextgengrafana/ticdc_new_arch_next_gen.json b/metrics/nextgengrafana/ticdc_new_arch_next_gen.json index f920ac2ccd..40b69b9b79 100644 --- a/metrics/nextgengrafana/ticdc_new_arch_next_gen.json +++ b/metrics/nextgengrafana/ticdc_new_arch_next_gen.json @@ -19166,6 +19166,110 @@ "align": false, "alignLevel": null } + }, + { + "aliasColors": {}, + "bars": false, + "dashLength": 10, + "dashes": false, + "datasource": "${DS_TEST-CLUSTER}", + "description": "Average and p99 throttle time reported by each broker.", + "fieldConfig": { + "defaults": { + "links": [] + }, + "overrides": [] + }, + "fill": 1, + "fillGradient": 0, + "gridPos": { + "h": 7, + "w": 12, + "x": 12, + "y": 67 + }, + "hiddenSeries": false, + "id": 62108, + "legend": { + "alignAsTable": true, + "avg": false, + "current": true, + "max": true, + "min": false, + "rightSide": false, + "show": true, + "total": false, + "values": true + }, + "lines": true, + "linewidth": 1, + "links": [], + "nullPointMode": "null", + "options": { + "alertThreshold": true + }, + "paceLength": 10, + "percentage": false, + "pluginVersion": "7.5.17", + "pointradius": 2, + "points": false, + "renderer": "flot", + "seriesOverrides": [], + "spaceLength": 10, + "stack": false, + "steppedLine": false, + "targets": [ + { + "exemplar": true, + "expr": "sum(ticdc_sink_kafka_producer_throttle_time{k8s_cluster=\"$k8s_cluster\", sharedpool_id=\"$tidb_cluster\", keyspace_name=~\"$keyspace_name\",changefeed=~\"$changefeed\", instance=~\"$ticdc_instance\"}) by (keyspace_name,changefeed, instance, broker, type)", + "format": "time_series", + "interval": "", + "intervalFactor": 1, + "legendFormat": "{{keyspace_name}}-{{changefeed}}-{{instance}}-{{broker}}-{{type}}", + "refId": "A" + } + ], + "thresholds": [], + "timeFrom": null, + "timeRegions": [], + "timeShift": null, + "title": "Kafka Throttle Time", + "tooltip": { + "shared": true, + "sort": 0, + "value_type": "individual" + }, + "type": "graph", + "xaxis": { + "buckets": null, + "mode": "time", + "name": null, + "show": true, + "values": [] + }, + "yaxes": [ + { + "decimals": 1, + "format": "ms", + "label": null, + "logBase": 1, + "max": null, + "min": "0", + "show": true + }, + { + "format": "short", + "label": null, + "logBase": 1, + "max": null, + "min": null, + "show": false + } + ], + "yaxis": { + "align": false, + "alignLevel": null + } } ], "title": "Sink - MQ Sink", @@ -28342,5 +28446,5 @@ "timezone": "browser", "title": "${DS_TEST-CLUSTER}-TiCDC-New-Arch", "uid": "YiGL8hBZ0aac", - "version": 41 + "version": 42 } diff --git a/metrics/nextgengrafana/ticdc_new_arch_with_keyspace_name.json b/metrics/nextgengrafana/ticdc_new_arch_with_keyspace_name.json index f71e45145a..efddefb111 100644 --- a/metrics/nextgengrafana/ticdc_new_arch_with_keyspace_name.json +++ b/metrics/nextgengrafana/ticdc_new_arch_with_keyspace_name.json @@ -7691,6 +7691,110 @@ "align": false, "alignLevel": null } + }, + { + "aliasColors": {}, + "bars": false, + "dashLength": 10, + "dashes": false, + "datasource": "${DS_TEST-CLUSTER}", + "description": "Average and p99 throttle time reported by each broker.", + "fieldConfig": { + "defaults": { + "links": [] + }, + "overrides": [] + }, + "fill": 1, + "fillGradient": 0, + "gridPos": { + "h": 7, + "w": 12, + "x": 12, + "y": 67 + }, + "hiddenSeries": false, + "id": 62108, + "legend": { + "alignAsTable": true, + "avg": false, + "current": true, + "max": true, + "min": false, + "rightSide": false, + "show": true, + "total": false, + "values": true + }, + "lines": true, + "linewidth": 1, + "links": [], + "nullPointMode": "null", + "options": { + "alertThreshold": true + }, + "paceLength": 10, + "percentage": false, + "pluginVersion": "7.5.17", + "pointradius": 2, + "points": false, + "renderer": "flot", + "seriesOverrides": [], + "spaceLength": 10, + "stack": false, + "steppedLine": false, + "targets": [ + { + "exemplar": true, + "expr": "sum(ticdc_sink_kafka_producer_throttle_time{k8s_cluster=\"$k8s_cluster\", tidb_cluster=\"$tidb_cluster\", keyspace_name=~\"$keyspace_name\",changefeed=~\"$changefeed\", instance=~\"$ticdc_instance\"}) by (keyspace_name,changefeed, instance, broker, type)", + "format": "time_series", + "interval": "", + "intervalFactor": 1, + "legendFormat": "{{keyspace_name}}-{{changefeed}}-{{instance}}-{{broker}}-{{type}}", + "refId": "A" + } + ], + "thresholds": [], + "timeFrom": null, + "timeRegions": [], + "timeShift": null, + "title": "Kafka Throttle Time", + "tooltip": { + "shared": true, + "sort": 0, + "value_type": "individual" + }, + "type": "graph", + "xaxis": { + "buckets": null, + "mode": "time", + "name": null, + "show": true, + "values": [] + }, + "yaxes": [ + { + "decimals": 1, + "format": "ms", + "label": null, + "logBase": 1, + "max": null, + "min": "0", + "show": true + }, + { + "format": "short", + "label": null, + "logBase": 1, + "max": null, + "min": null, + "show": false + } + ], + "yaxis": { + "align": false, + "alignLevel": null + } } ], "title": "Sink - MQ Sink", @@ -11749,5 +11853,5 @@ "timezone": "browser", "title": "${DS_TEST-CLUSTER}-TiCDC-New-Arch-KeyspaceName", "uid": "lGT5hED6vqTn", - "version": 41 + "version": 42 } diff --git a/pkg/sink/kafka/metrics.go b/pkg/sink/kafka/metrics.go index d3f88055e2..684457e420 100644 --- a/pkg/sink/kafka/metrics.go +++ b/pkg/sink/kafka/metrics.go @@ -69,6 +69,13 @@ var ( Name: "kafka_producer_records_per_request", Help: "The number of records per request for all topics.", }, []string{"namespace", "changefeed", "type"}) + throttleTimeGauge = prometheus.NewGaugeVec( + prometheus.GaugeOpts{ + Namespace: "ticdc", + Subsystem: "sink", + Name: "kafka_producer_throttle_time", + Help: "Kafka broker throttle time in milliseconds.", + }, []string{"namespace", "changefeed", "broker", "type"}) // Meter mark by 1 once a response received. responseRateGauge = prometheus.NewGaugeVec( @@ -84,6 +91,7 @@ var ( func InitMetrics(registry *prometheus.Registry) { registry.MustRegister(compressionRatioGauge) registry.MustRegister(recordsPerRequestGauge) + registry.MustRegister(throttleTimeGauge) registry.MustRegister(OutgoingByteRateGauge) registry.MustRegister(RequestRateGauge) registry.MustRegister(RequestLatencyGauge) diff --git a/pkg/sink/kafka/metrics_collector.go b/pkg/sink/kafka/metrics_collector.go index 3e54c7d82f..03bb12a89a 100644 --- a/pkg/sink/kafka/metrics_collector.go +++ b/pkg/sink/kafka/metrics_collector.go @@ -50,6 +50,7 @@ const ( requestLatencyInMsMetricNamePrefix = "request-latency-in-ms-for-broker-" requestsInFlightMetricNamePrefix = "requests-in-flight-for-broker-" responseRateMetricNamePrefix = "response-rate-for-broker-" + throttleTimeMetricNamePrefix = "throttle-time-in-ms-for-broker-" p99 = "p99" avg = "avg" @@ -168,6 +169,18 @@ func (m *saramaMetricsCollector) collectBrokerMetrics() { WithLabelValues(keyspace, changefeedID, brokerID). Set(meter.Snapshot().Rate1()) } + + throttleTimeMetric := m.registry.Get(getBrokerMetricName( + throttleTimeMetricNamePrefix, brokerID)) + if histogram, ok := throttleTimeMetric.(metrics.Histogram); ok { + snapshot := histogram.Snapshot() + throttleTimeGauge. + WithLabelValues(keyspace, changefeedID, brokerID, avg). + Set(snapshot.Mean()) + throttleTimeGauge. + WithLabelValues(keyspace, changefeedID, brokerID, p99). + Set(snapshot.Percentile(0.99)) + } } } @@ -204,6 +217,10 @@ func (m *saramaMetricsCollector) cleanupBrokerMetrics() { DeleteLabelValues(keyspace, changefeedID, brokerID) responseRateGauge. DeleteLabelValues(keyspace, changefeedID, brokerID) + throttleTimeGauge. + DeleteLabelValues(keyspace, changefeedID, brokerID, avg) + throttleTimeGauge. + DeleteLabelValues(keyspace, changefeedID, brokerID, p99) } } From e9c917ff40b8b0a26d118d2640c3bd1cd42f8555 Mon Sep 17 00:00:00 2001 From: 3AceShowHand Date: Mon, 17 Aug 2026 16:44:14 +0800 Subject: [PATCH 2/2] add version to the kafka --- metrics/grafana/ticdc_new_arch.json | 2 +- metrics/nextgengrafana/ticdc_new_arch_next_gen.json | 2 +- metrics/nextgengrafana/ticdc_new_arch_with_keyspace_name.json | 2 +- pkg/sink/kafka/metrics.go | 2 +- pkg/sink/kafka/metrics_collector.go | 4 ++-- 5 files changed, 6 insertions(+), 6 deletions(-) diff --git a/metrics/grafana/ticdc_new_arch.json b/metrics/grafana/ticdc_new_arch.json index 8dda9ef46a..488f76d123 100644 --- a/metrics/grafana/ticdc_new_arch.json +++ b/metrics/grafana/ticdc_new_arch.json @@ -19250,7 +19250,7 @@ "yaxes": [ { "decimals": 1, - "format": "ms", + "format": "s", "label": null, "logBase": 1, "max": null, diff --git a/metrics/nextgengrafana/ticdc_new_arch_next_gen.json b/metrics/nextgengrafana/ticdc_new_arch_next_gen.json index 40b69b9b79..d7e7555d3b 100644 --- a/metrics/nextgengrafana/ticdc_new_arch_next_gen.json +++ b/metrics/nextgengrafana/ticdc_new_arch_next_gen.json @@ -19250,7 +19250,7 @@ "yaxes": [ { "decimals": 1, - "format": "ms", + "format": "s", "label": null, "logBase": 1, "max": null, diff --git a/metrics/nextgengrafana/ticdc_new_arch_with_keyspace_name.json b/metrics/nextgengrafana/ticdc_new_arch_with_keyspace_name.json index efddefb111..d17672f4b1 100644 --- a/metrics/nextgengrafana/ticdc_new_arch_with_keyspace_name.json +++ b/metrics/nextgengrafana/ticdc_new_arch_with_keyspace_name.json @@ -7775,7 +7775,7 @@ "yaxes": [ { "decimals": 1, - "format": "ms", + "format": "s", "label": null, "logBase": 1, "max": null, diff --git a/pkg/sink/kafka/metrics.go b/pkg/sink/kafka/metrics.go index 684457e420..fc1ebf7594 100644 --- a/pkg/sink/kafka/metrics.go +++ b/pkg/sink/kafka/metrics.go @@ -74,7 +74,7 @@ var ( Namespace: "ticdc", Subsystem: "sink", Name: "kafka_producer_throttle_time", - Help: "Kafka broker throttle time in milliseconds.", + Help: "Kafka broker throttle time in seconds.", }, []string{"namespace", "changefeed", "broker", "type"}) // Meter mark by 1 once a response received. diff --git a/pkg/sink/kafka/metrics_collector.go b/pkg/sink/kafka/metrics_collector.go index 03bb12a89a..674afa7a7f 100644 --- a/pkg/sink/kafka/metrics_collector.go +++ b/pkg/sink/kafka/metrics_collector.go @@ -176,10 +176,10 @@ func (m *saramaMetricsCollector) collectBrokerMetrics() { snapshot := histogram.Snapshot() throttleTimeGauge. WithLabelValues(keyspace, changefeedID, brokerID, avg). - Set(snapshot.Mean()) + Set(snapshot.Mean() / 1000) throttleTimeGauge. WithLabelValues(keyspace, changefeedID, brokerID, p99). - Set(snapshot.Percentile(0.99)) + Set(snapshot.Percentile(0.99) / 1000) } } }