Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
43 changes: 43 additions & 0 deletions docs/content/metrics/domains/stmt.schema/_index.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,43 @@
---
title: "stmt.schema"
---

The `stmt.schema` domain sums statement counters from MySQL 8.0 `performance_schema.events_statements_summary_by_digest` by `SCHEMA_NAME`. It is opt-in and does not change the default plan.

## Usage

```yaml
level:
collect:
stmt.schema:
metrics:
- count_star
- sum_timer_wait
- sum_errors
- sum_rows_examined
options:
include: app,other_app
```

The default metrics are `count_star`, `sum_timer_wait`, `sum_errors`, `sum_rows_examined`, and `sum_rows_sent`. Other supported metrics are `sum_lock_time`, `sum_warnings`, `sum_rows_affected`, `sum_created_tmp_tables`, `sum_created_tmp_disk_tables`, `sum_select_scan`, `sum_select_full_join`, `sum_no_index_used`, and `sum_no_good_index_used`.

All values are cumulative counters. Timer values are in picoseconds. Blip's delta sink converts cumulative values to changes between samples. The `db` group contains the statement schema, which does not necessarily identify every table touched by a cross-schema query.

Rows with a null schema or digest are excluded. This includes the digest catch-all row, so the series can undercount when `performance_schema_digests_size` is exhausted. Monitor the catch-all row separately when completeness matters.

## Options

|Option|Default|Description|
|------|-------|-----------|
|`include`||Comma-separated schema names to include; overrides `exclude`|
|`exclude`|`mysql,information_schema,performance_schema,sys`|Comma-separated schema names to exclude|

## Group Keys

|Key|Value|
|---|-----|
|`db`|Statement schema name|

## MySQL Config

Enable Performance Schema and its `statements_digest` consumer. The digest table has a fixed size chosen at startup. A full table sends new digests to a null catch-all row, which this collector cannot assign to a schema.
15 changes: 13 additions & 2 deletions docs/content/metrics/domains/wait.io.table/_index.md
Original file line number Diff line number Diff line change
Expand Up @@ -70,14 +70,24 @@ level:
All Blip metric names are lowercase when reported.
{{< /hint >}}

Metrics are [grouped](#group-keys) by database and table.
By default, metrics are [grouped](#group-keys) by database and table. Set `group-by: schema` to report one series per database instead.

## Derived Metrics

None.

## Options

### `group-by`

|Value|Default|Description|
|-----|-------|-----------|
|table|&check;|One series per table, grouped by `db` and `tbl`|
|schema||One series per schema, grouped by `db`|
|both||Table and schema series|

Schema rollups sum counts and total wait times across tables. Minimum and maximum times span the tables with events; averages are weighted by the event count. These are table handler waits, not physical disk I/O or complete query latency.

### `all`

|Value|Default|Description|
Expand Down Expand Up @@ -127,7 +137,8 @@ Normally, truncating a table is nearly instantaneous, but metadata locks can blo

|Key|Value|
|---|---|
|`db`, `tbl`|Database and table name|
|`db`, `tbl`|Database and table name for table series|
|`db`|Database name for schema series|

## Meta

Expand Down
3 changes: 2 additions & 1 deletion docs/content/metrics/quick-ref.md
Original file line number Diff line number Diff line change
Expand Up @@ -3,7 +3,7 @@ weight: 100
---

Following are _all_ Blip domains and the metrics collected in each.
Only domains with a Blip version are collected.
Only domains with a Blip version or `Unreleased` are collected.
The rest are reserved for future use.

|Domain|Metrics|Blip Version|
Expand Down Expand Up @@ -65,6 +65,7 @@ The rest are reserved for future use.
|status.user|Status by user||
|stmt|Statements||
|[`stmt.current`](domains#stmtcurrent)|Current statements|v1.0.0|
|[`stmt.schema`](domains#stmtschema)|Statement counters by schema (`events_statements_summary_by_digest`)|Unreleased|
|stmt.history|Historical statements||
|thd|Threads||
|[`tls`](domains#tls)|TLS (SSL) status and configuration|v1.0.0|
Expand Down
4 changes: 4 additions & 0 deletions metrics/factory.go
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,7 @@ import (
sizetable "github.com/cashapp/blip/v2/metrics/size.table"
statusglobal "github.com/cashapp/blip/v2/metrics/status.global"
"github.com/cashapp/blip/v2/metrics/stmt.current"
stmtschema "github.com/cashapp/blip/v2/metrics/stmt.schema"
"github.com/cashapp/blip/v2/metrics/tls"
"github.com/cashapp/blip/v2/metrics/trx"
varglobal "github.com/cashapp/blip/v2/metrics/var.global"
Expand Down Expand Up @@ -418,6 +419,8 @@ func (f *factory) Make(domain string, args blip.CollectorFactoryArgs) (blip.Coll
return statusglobal.NewGlobal(args.DB), nil
case "stmt.current":
return stmt.NewCurrent(args.DB), nil
case "stmt.schema":
return stmtschema.NewSchema(args.DB), nil
case "tls":
return tls.NewTLS(args.DB), nil
case "trx":
Expand Down Expand Up @@ -451,6 +454,7 @@ var builtinCollectors = []string{
"size.table",
"status.global",
"stmt.current",
"stmt.schema",
"trx",
"tls",
"var.global",
Expand Down
183 changes: 183 additions & 0 deletions metrics/stmt.schema/schema.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,183 @@
// Copyright 2026 Block, Inc.

package stmtschema

import (
"context"
"database/sql"
"fmt"
"strings"

"github.com/cashapp/blip/v2"
"github.com/cashapp/blip/v2/sqlutil"
)

const (
DOMAIN = "stmt.schema"

OPT_INCLUDE = "include"
OPT_EXCLUDE = "exclude"

defaultExclude = "mysql,information_schema,performance_schema,sys"
)

var availableMetrics = []string{
"count_star",
"sum_timer_wait",
"sum_lock_time",
"sum_errors",
"sum_warnings",
"sum_rows_affected",
"sum_rows_sent",
"sum_rows_examined",
"sum_created_tmp_tables",
"sum_created_tmp_disk_tables",
"sum_select_scan",
"sum_select_full_join",
"sum_no_index_used",
"sum_no_good_index_used",
}

var defaultMetrics = []string{
"count_star", "sum_timer_wait", "sum_errors", "sum_rows_examined", "sum_rows_sent",
}

type config struct {
query string
params []interface{}
metrics []string
}

// Schema collects cumulative statement counters by the statement's schema.
type Schema struct {
db *sql.DB
atLevel map[string]config
}

var _ blip.Collector = &Schema{}

func NewSchema(db *sql.DB) *Schema {
return &Schema{db: db, atLevel: map[string]config{}}
}

func (c *Schema) Domain() string { return DOMAIN }

func (c *Schema) Help() blip.CollectorHelp {
metrics := make([]blip.CollectorMetric, 0, len(availableMetrics))
for _, name := range availableMetrics {
metrics = append(metrics, blip.CollectorMetric{Name: name, Type: blip.CUMULATIVE_COUNTER})
}
return blip.CollectorHelp{
Domain: DOMAIN,
Description: "Statement counters grouped by Performance Schema statement schema",
Options: map[string]blip.CollectorHelpOption{
OPT_INCLUDE: {Name: OPT_INCLUDE, Desc: "Comma-separated schema names to include"},
OPT_EXCLUDE: {Name: OPT_EXCLUDE, Desc: "Comma-separated schema names to exclude when include is unset", Default: defaultExclude},
},
Groups: []blip.CollectorKeyValue{{Key: "db", Value: "statement schema name"}},
Metrics: metrics,
}
}

func (c *Schema) Prepare(ctx context.Context, plan blip.Plan) (func(), error) {
for _, level := range plan.Levels {
dom, ok := level.Collect[DOMAIN]
if !ok {
continue
}
query, params, names, err := SummaryQuery(dom.Options, dom.Metrics)
if err != nil {
return nil, err
}
c.atLevel[level.Name] = config{query: query, params: params, metrics: names}
}
return nil, nil
}

func (c *Schema) Collect(ctx context.Context, levelName string) ([]blip.MetricValue, error) {
config, ok := c.atLevel[levelName]
if !ok {
return nil, nil
}
rows, err := c.db.QueryContext(ctx, config.query, config.params...)
if err != nil {
return nil, err
}
defer rows.Close()

var values []blip.MetricValue
for rows.Next() {
var schema string
counts := make([]float64, len(config.metrics))
dest := make([]interface{}, len(counts)+1)
dest[0] = &schema
for i := range counts {
dest[i+1] = &counts[i]
}
if err := rows.Scan(dest...); err != nil {
return nil, err
}
for i, name := range config.metrics {
values = append(values, blip.MetricValue{
Name: name, Value: counts[i], Type: blip.CUMULATIVE_COUNTER,
Group: map[string]string{"db": schema},
})
}
}
return values, rows.Err()
}

// SummaryQuery selects only known additive counters. The digest catch-all row
// has no schema, so it cannot be assigned to a schema-level series.
func SummaryQuery(opts map[string]string, requested []string) (string, []interface{}, []string, error) {
if len(requested) == 0 {
requested = defaultMetrics
}
allowed := make(map[string]bool, len(availableMetrics))
for _, metric := range availableMetrics {
allowed[metric] = true
}
names := make([]string, 0, len(requested))
seen := map[string]bool{}
for _, metric := range requested {
name := strings.ToLower(metric)
if !allowed[name] {
return "", nil, nil, fmt.Errorf("invalid %s metric: %s", DOMAIN, metric)
}
if !seen[name] {
names = append(names, name)
seen[name] = true
}
}

columns := make([]string, 0, len(names)+1)
columns = append(columns, "SCHEMA_NAME")
for _, name := range names {
columns = append(columns, fmt.Sprintf("SUM(%s) AS %s", name, name))
}
query := "SELECT " + strings.Join(columns, ", ") +
" FROM performance_schema.events_statements_summary_by_digest WHERE SCHEMA_NAME IS NOT NULL AND DIGEST IS NOT NULL"

var schemas []string
include := false
if opts[OPT_INCLUDE] != "" {
schemas = strings.Split(opts[OPT_INCLUDE], ",")
include = true
} else {
exclude := opts[OPT_EXCLUDE]
if exclude == "" {
exclude = defaultExclude
}
schemas = strings.Split(exclude, ",")
}
params := make([]interface{}, len(schemas))
for i, schema := range schemas {
params[i] = strings.TrimSpace(schema)
}
if include {
query += " AND SCHEMA_NAME IN (" + sqlutil.PlaceholderList(len(schemas)) + ")"
} else {
query += " AND SCHEMA_NAME NOT IN (" + sqlutil.PlaceholderList(len(schemas)) + ")"
}
return query + " GROUP BY SCHEMA_NAME", params, names, nil
}
81 changes: 81 additions & 0 deletions metrics/stmt.schema/schema_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,81 @@
// Copyright 2026 Block, Inc.

package stmtschema

import (
"context"
"strings"
"testing"

"github.com/cashapp/blip/v2"
"github.com/cashapp/blip/v2/test"
)

func TestSummaryQuery(t *testing.T) {
query, params, names, err := SummaryQuery(map[string]string{
OPT_INCLUDE: "app, other",
OPT_EXCLUDE: "ignored",
}, []string{"COUNT_STAR", "sum_errors", "count_star"})
if err != nil {
t.Fatal(err)
}
if want := "SELECT SCHEMA_NAME, SUM(count_star) AS count_star, SUM(sum_errors) AS sum_errors FROM performance_schema.events_statements_summary_by_digest WHERE SCHEMA_NAME IS NOT NULL AND DIGEST IS NOT NULL AND SCHEMA_NAME IN (?, ?) GROUP BY SCHEMA_NAME"; query != want {
t.Errorf("query:\n%s\nwant:\n%s", query, want)
}
if len(names) != 2 || names[0] != "count_star" || names[1] != "sum_errors" {
t.Errorf("names = %v", names)
}
if len(params) != 2 || params[0] != "app" || params[1] != "other" {
t.Errorf("params = %v", params)
}

query, _, names, err = SummaryQuery(nil, nil)
if err != nil || len(names) == 0 || !strings.Contains(query, "SCHEMA_NAME NOT IN (?, ?, ?, ?)") {
t.Errorf("default query = %q, names = %v, err = %v", query, names, err)
}
}

func TestCollectMySQL80(t *testing.T) {
_, db, err := test.Connection("mysql80")
if err != nil {
t.Skip("mysql80 not running")
}
defer db.Close()

ctx := context.Background()
conn, err := db.Conn(ctx)
if err != nil {
t.Fatal(err)
}
defer conn.Close()
if _, err := conn.ExecContext(ctx, "USE mysql"); err != nil {
t.Fatal(err)
}
if _, err := conn.ExecContext(ctx, "SELECT 42"); err != nil {
t.Fatal(err)
}

c := NewSchema(db)
plan := blip.Plan{Levels: map[string]blip.Level{
"test": {Name: "test", Collect: map[string]blip.Domain{
DOMAIN: {Metrics: []string{"count_star", "sum_rows_sent"}, Options: map[string]string{OPT_INCLUDE: "mysql"}},
}},
}}
if _, err := c.Prepare(ctx, plan); err != nil {
t.Fatal(err)
}
values, err := c.Collect(ctx, "test")
if err != nil {
t.Fatal(err)
}
if len(values) != 2 || values[0].Group["db"] != "mysql" || values[0].Value < 1 {
t.Fatalf("unexpected schema metrics: %+v", values)
}
}

func TestSummaryQueryRejectsUnknownMetric(t *testing.T) {
_, _, _, err := SummaryQuery(nil, []string{"avg_timer_wait"})
if err == nil {
t.Fatal("expected unknown metric error")
}
}
Loading
Loading