From 46d3479986c6b46a4316bd8a7684ba708f06a617 Mon Sep 17 00:00:00 2001 From: Ian Oberst Date: Thu, 24 Sep 2026 10:46:29 -0700 Subject: [PATCH 1/2] Add schema-level MySQL statement and I/O metrics --- .../metrics/domains/stmt.schema/_index.md | 43 ++++ .../metrics/domains/wait.io.table/_index.md | 15 +- docs/content/metrics/quick-ref.md | 3 +- metrics/factory.go | 4 + metrics/stmt.schema/schema.go | 183 ++++++++++++++++++ metrics/stmt.schema/schema_test.go | 81 ++++++++ metrics/wait.io.table/query.go | 36 ++++ metrics/wait.io.table/query_test.go | 60 ++++++ metrics/wait.io.table/table.go | 86 +++++--- 9 files changed, 482 insertions(+), 29 deletions(-) create mode 100644 docs/content/metrics/domains/stmt.schema/_index.md create mode 100644 metrics/stmt.schema/schema.go create mode 100644 metrics/stmt.schema/schema_test.go diff --git a/docs/content/metrics/domains/stmt.schema/_index.md b/docs/content/metrics/domains/stmt.schema/_index.md new file mode 100644 index 00000000..7748cc05 --- /dev/null +++ b/docs/content/metrics/domains/stmt.schema/_index.md @@ -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. diff --git a/docs/content/metrics/domains/wait.io.table/_index.md b/docs/content/metrics/domains/wait.io.table/_index.md index ce804472..21ae79bf 100644 --- a/docs/content/metrics/domains/wait.io.table/_index.md +++ b/docs/content/metrics/domains/wait.io.table/_index.md @@ -70,7 +70,7 @@ 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 @@ -78,6 +78,16 @@ None. ## Options +### `group-by` + +|Value|Default|Description| +|-----|-------|-----------| +|table|✓|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| @@ -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 diff --git a/docs/content/metrics/quick-ref.md b/docs/content/metrics/quick-ref.md index 4fab1e5c..c3f8ab41 100644 --- a/docs/content/metrics/quick-ref.md +++ b/docs/content/metrics/quick-ref.md @@ -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| @@ -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| diff --git a/metrics/factory.go b/metrics/factory.go index b50b4dcb..5df71de2 100644 --- a/metrics/factory.go +++ b/metrics/factory.go @@ -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" @@ -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": @@ -451,6 +454,7 @@ var builtinCollectors = []string{ "size.table", "status.global", "stmt.current", + "stmt.schema", "trx", "tls", "var.global", diff --git a/metrics/stmt.schema/schema.go b/metrics/stmt.schema/schema.go new file mode 100644 index 00000000..98b1cce0 --- /dev/null +++ b/metrics/stmt.schema/schema.go @@ -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 +} diff --git a/metrics/stmt.schema/schema_test.go b/metrics/stmt.schema/schema_test.go new file mode 100644 index 00000000..7e5e065d --- /dev/null +++ b/metrics/stmt.schema/schema_test.go @@ -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") + } +} diff --git a/metrics/wait.io.table/query.go b/metrics/wait.io.table/query.go index 7147aac7..91ad7e88 100644 --- a/metrics/wait.io.table/query.go +++ b/metrics/wait.io.table/query.go @@ -21,6 +21,42 @@ func TableIoWaitQuery(set map[string]string, metrics []string) (string, []interf return query + where, params } +// TableIoWaitSchemaQuery rolls up table handler counters by the schema that owns +// each table. The count columns provide the weights for average wait times. +func TableIoWaitSchemaQuery(set map[string]string, metrics []string) (string, []interface{}) { + columns := setColumns(set, metrics)[2:] + selectColumns := []string{"OBJECT_SCHEMA", "'' AS OBJECT_NAME"} + for _, column := range columns { + switch { + case strings.HasPrefix(column, "count_"), strings.HasPrefix(column, "sum_timer_"): + selectColumns = append(selectColumns, fmt.Sprintf("SUM(%s) AS %s", column, column)) + case strings.HasPrefix(column, "min_timer_"): + selectColumns = append(selectColumns, fmt.Sprintf("COALESCE(MIN(CASE WHEN count_%s > 0 THEN %s END), 0) AS %s", countSuffix(column), column, column)) + case strings.HasPrefix(column, "avg_timer_"): + operation := strings.TrimPrefix(column, "avg_timer_") + selectColumns = append(selectColumns, fmt.Sprintf("COALESCE(CAST(SUM(sum_timer_%s) / NULLIF(SUM(count_%s), 0) AS UNSIGNED), 0) AS %s", operation, countSuffix(column), column)) + case strings.HasPrefix(column, "max_timer_"): + selectColumns = append(selectColumns, fmt.Sprintf("MAX(%s) AS %s", column, column)) + } + } + var where string + var params []interface{} + if include := set[OPT_INCLUDE]; include != "" { + where, params = setWhere(strings.Split(include, ","), true) + } else { + where, params = setWhere(strings.Split(set[OPT_EXCLUDE], ","), false) + } + return fmt.Sprintf("SELECT %s FROM performance_schema.table_io_waits_summary_by_table%s GROUP BY OBJECT_SCHEMA", strings.Join(selectColumns, ", "), where), params +} + +func countSuffix(timerColumn string) string { + operation := timerColumn[strings.LastIndex(timerColumn, "_timer_")+len("_timer_"):] + if operation == "wait" { + return "star" + } + return operation +} + func setColumns(set map[string]string, metrics []string) []string { columns := []string{"OBJECT_SCHEMA", "OBJECT_NAME"} diff --git a/metrics/wait.io.table/query_test.go b/metrics/wait.io.table/query_test.go index 7c49c355..44edf4ec 100644 --- a/metrics/wait.io.table/query_test.go +++ b/metrics/wait.io.table/query_test.go @@ -3,9 +3,12 @@ package waitiotable_test import ( + "context" "testing" + "github.com/cashapp/blip/v2" waitiotable "github.com/cashapp/blip/v2/metrics/wait.io.table" + "github.com/cashapp/blip/v2/test" "github.com/go-test/deep" ) @@ -64,3 +67,60 @@ func TestTableIoQuery(t *testing.T) { t.Error(diff) } } + +func TestCollectSchemaRollupMySQL80(t *testing.T) { + _, db, err := test.Connection("mysql80") + if err != nil { + t.Skip("mysql80 not running") + } + defer db.Close() + + c := waitiotable.NewTable(db) + plan := blip.Plan{Levels: map[string]blip.Level{ + "test": {Name: "test", Collect: map[string]blip.Domain{ + waitiotable.DOMAIN: { + Metrics: []string{"count_star", "min_timer_wait", "avg_timer_wait", "max_timer_wait"}, + Options: map[string]string{ + waitiotable.OPT_INCLUDE: "mysql.*", + waitiotable.OPT_GROUP_BY: "both", + waitiotable.OPT_TRUNCATE_TABLE: "no", + }, + }, + }}, + }} + ctx := context.Background() + if _, err := c.Prepare(ctx, plan); err != nil { + t.Fatal(err) + } + values, err := c.Collect(ctx, "test") + if err != nil { + t.Fatal(err) + } + var table, schema bool + for _, value := range values { + if value.Group["db"] != "mysql" { + t.Fatalf("unexpected group: %+v", value.Group) + } + if _, ok := value.Group["tbl"]; ok { + table = true + } else { + schema = true + } + } + if !table || !schema { + t.Fatalf("expected both table and schema metrics, got %d values", len(values)) + } +} + +func TestTableIoSchemaQuery(t *testing.T) { + opts := map[string]string{waitiotable.OPT_INCLUDE: "app.*,other.orders"} + metrics := []string{"count_fetch", "sum_timer_fetch", "min_timer_fetch", "avg_timer_fetch", "max_timer_fetch"} + got, params := waitiotable.TableIoWaitSchemaQuery(opts, metrics) + expect := "SELECT OBJECT_SCHEMA, '' AS OBJECT_NAME, SUM(count_fetch) AS count_fetch, SUM(sum_timer_fetch) AS sum_timer_fetch, COALESCE(MIN(CASE WHEN count_fetch > 0 THEN min_timer_fetch END), 0) AS min_timer_fetch, COALESCE(CAST(SUM(sum_timer_fetch) / NULLIF(SUM(count_fetch), 0) AS UNSIGNED), 0) AS avg_timer_fetch, MAX(max_timer_fetch) AS max_timer_fetch FROM performance_schema.table_io_waits_summary_by_table WHERE (OBJECT_SCHEMA = ?) OR (OBJECT_SCHEMA = ? AND OBJECT_NAME = ?) GROUP BY OBJECT_SCHEMA" + if got != expect { + t.Errorf("got:\n%s\nexpect:\n%s", got, expect) + } + if diff := deep.Equal(params, []interface{}{"app", "other", "orders"}); diff != nil { + t.Error(diff) + } +} diff --git a/metrics/wait.io.table/table.go b/metrics/wait.io.table/table.go index 02ab0721..40655cff 100644 --- a/metrics/wait.io.table/table.go +++ b/metrics/wait.io.table/table.go @@ -21,6 +21,7 @@ const ( OPT_TRUNCATE_TABLE = "truncate-table" OPT_TRUNCATE_TIMEOUT = "truncate-timeout" OPT_ALL = "all" + OPT_GROUP_BY = "group-by" OPT_EXCLUDE_DEFAULT = "mysql.*,information_schema.*,performance_schema.*,sys.*" @@ -82,6 +83,9 @@ func init() { type tableOptions struct { query string params []interface{} + schemaQuery string + schemaParams []interface{} + groupBy string truncate bool truncateTimeout time.Duration stop bool @@ -151,6 +155,16 @@ func (t *Table) Help() blip.CollectorHelp { "no": "Specified metrics", }, }, + OPT_GROUP_BY: { + Name: OPT_GROUP_BY, + Desc: "Group table I/O by table, owning schema, or both", + Default: "table", + Values: map[string]string{ + "table": "One series per table", + "schema": "One series per schema", + "both": "Table and schema series", + }, + }, }, Groups: []blip.CollectorKeyValue{ {Key: "db", Value: "the database name for the corresponding table io, or empty string for all dbs"}, @@ -184,6 +198,10 @@ LEVEL: } o.query, o.params = TableIoWaitQuery(dom.Options, dom.Metrics) + o.groupBy = dom.Options[OPT_GROUP_BY] + if o.groupBy == "schema" || o.groupBy == "both" { + o.schemaQuery, o.schemaParams = TableIoWaitSchemaQuery(dom.Options, dom.Metrics) + } if truncate, ok := dom.Options[OPT_TRUNCATE_TABLE]; ok && truncate == "no" { o.truncate = false @@ -234,7 +252,41 @@ func (t *Table) Collect(ctx context.Context, levelName string) ([]blip.MetricVal return nil, nil } - rows, err := t.db.QueryContext(ctx, o.query, o.params...) + var metrics []blip.MetricValue + if o.groupBy != "schema" { + values, err := t.collectQuery(ctx, o.query, o.params, o.metricType, false) + if err != nil { + return nil, err + } + metrics = append(metrics, values...) + } + if o.schemaQuery != "" { + values, err := t.collectQuery(ctx, o.schemaQuery, o.schemaParams, o.metricType, true) + if err != nil { + return nil, err + } + metrics = append(metrics, values...) + } + + if o.truncate { + conn, err := t.db.Conn(ctx) + if err == nil { + defer conn.Close() + _, err = conn.ExecContext(ctx, o.lockWaitQuery) + if err == nil { + trCtx, cancelFn := context.WithTimeout(ctx, o.truncateTimeout) + defer cancelFn() + _, err = conn.ExecContext(trCtx, TRUNCATE_QUERY) + } + } + return o.truncateErrPolicy.TruncateError(err, &o.stop, metrics) + } + + return metrics, nil +} + +func (t *Table) collectQuery(ctx context.Context, query string, params []interface{}, metricType byte, schema bool) ([]blip.MetricValue, error) { + rows, err := t.db.QueryContext(ctx, query, params...) if err != nil { return nil, err } @@ -269,10 +321,14 @@ func (t *Table) Collect(ctx context.Context, levelName string) ([]blip.MetricVal tblName = *values[1].(*string) for i := 2; i < len(cols); i++ { + group := map[string]string{"db": dbName, "tbl": tblName} + if schema { + group = map[string]string{"db": dbName} + } m := blip.MetricValue{ Name: cols[i], - Type: o.metricType, - Group: map[string]string{"db": dbName, "tbl": tblName}, + Type: metricType, + Group: group, } m.Value = float64(*values[i].(*int64)) metrics = append(metrics, m) @@ -280,27 +336,5 @@ func (t *Table) Collect(ctx context.Context, levelName string) ([]blip.MetricVal } - if o.truncate { - conn, err := t.db.Conn(ctx) - if err == nil { - defer conn.Close() - - // Set `lock_wait_timeout` to prevent our query from being blocked for too long - // due to metadata locking. We treat a failure to set the lock wait timeout - // the same as a truncate timeout, as not setting creates a risk of having a thread - // hang for an extended period of time. - _, err = conn.ExecContext(ctx, o.lockWaitQuery) - if err == nil { - trCtx, cancelFn := context.WithTimeout(ctx, o.truncateTimeout) - defer cancelFn() - _, err = conn.ExecContext(trCtx, TRUNCATE_QUERY) - } - } - // Process any errors (or lack thereof) with the TruncateErrorPolicy as there is special handling - // for the metric values that need to be applied, even if there is not an error. See comments - // in `TruncateErrorPolicy` for more details. - return o.truncateErrPolicy.TruncateError(err, &o.stop, metrics) - } - - return metrics, err + return metrics, rows.Err() } From 6595bf9a85fc08b2397561213dbf81d3de8ca517 Mon Sep 17 00:00:00 2001 From: Ian Oberst Date: Thu, 24 Sep 2026 11:34:00 -0700 Subject: [PATCH 2/2] Fix schema I/O rollup overflow and test setup --- metrics/wait.io.table/query_test.go | 9 ++++- metrics/wait.io.table/table.go | 14 ++++++-- metrics/wait.io.table/table_overflow_test.go | 38 ++++++++++++++++++++ 3 files changed, 58 insertions(+), 3 deletions(-) create mode 100644 metrics/wait.io.table/table_overflow_test.go diff --git a/metrics/wait.io.table/query_test.go b/metrics/wait.io.table/query_test.go index 44edf4ec..6ef279e0 100644 --- a/metrics/wait.io.table/query_test.go +++ b/metrics/wait.io.table/query_test.go @@ -4,11 +4,14 @@ package waitiotable_test import ( "context" + "errors" + "net" "testing" "github.com/cashapp/blip/v2" waitiotable "github.com/cashapp/blip/v2/metrics/wait.io.table" "github.com/cashapp/blip/v2/test" + _ "github.com/go-sql-driver/mysql" "github.com/go-test/deep" ) @@ -71,7 +74,11 @@ func TestTableIoQuery(t *testing.T) { func TestCollectSchemaRollupMySQL80(t *testing.T) { _, db, err := test.Connection("mysql80") if err != nil { - t.Skip("mysql80 not running") + var netErr *net.OpError + if errors.As(err, &netErr) { + t.Skipf("mysql80 not running: %v", err) + } + t.Fatalf("connect to mysql80: %v", err) } defer db.Close() diff --git a/metrics/wait.io.table/table.go b/metrics/wait.io.table/table.go index 40655cff..929ee8be 100644 --- a/metrics/wait.io.table/table.go +++ b/metrics/wait.io.table/table.go @@ -309,7 +309,13 @@ func (t *Table) collectQuery(ctx context.Context, query string, params []interfa values[1] = new(string) for i := 2; i < len(cols); i++ { - values[i] = new(int64) + if schema { + // SUM of unsigned table counters can exceed int64, and schema + // metrics are represented as float64 after scanning. + values[i] = new(float64) + } else { + values[i] = new(int64) + } } for rows.Next() { @@ -330,7 +336,11 @@ func (t *Table) collectQuery(ctx context.Context, query string, params []interfa Type: metricType, Group: group, } - m.Value = float64(*values[i].(*int64)) + if schema { + m.Value = *values[i].(*float64) + } else { + m.Value = float64(*values[i].(*int64)) + } metrics = append(metrics, m) } diff --git a/metrics/wait.io.table/table_overflow_test.go b/metrics/wait.io.table/table_overflow_test.go new file mode 100644 index 00000000..aad66136 --- /dev/null +++ b/metrics/wait.io.table/table_overflow_test.go @@ -0,0 +1,38 @@ +// Copyright 2026 Block, Inc. + +package waitiotable + +import ( + "context" + "errors" + "math" + "net" + "testing" + + "github.com/cashapp/blip/v2" + "github.com/cashapp/blip/v2/test" + _ "github.com/go-sql-driver/mysql" +) + +func TestCollectSchemaLargeDecimalMySQL80(t *testing.T) { + _, db, err := test.Connection("mysql80") + if err != nil { + var netErr *net.OpError + if errors.As(err, &netErr) { + t.Skipf("mysql80 not running: %v", err) + } + t.Fatalf("connect to mysql80: %v", err) + } + defer db.Close() + + // MySQL returns SUM of unsigned counters as DECIMAL. This value is + // larger than MaxInt64, which the old scan destination could not hold. + query := "SELECT 'mysql' AS OBJECT_SCHEMA, '' AS OBJECT_NAME, CAST(18446744073709551616 AS DECIMAL(30, 0)) AS sum_timer_wait" + values, err := NewTable(db).collectQuery(context.Background(), query, nil, blip.CUMULATIVE_COUNTER, true) + if err != nil { + t.Fatal(err) + } + if len(values) != 1 || values[0].Value != math.Exp2(64) || values[0].Group["db"] != "mysql" { + t.Fatalf("unexpected schema metrics: %+v", values) + } +}