diff --git a/docs/reference/configuration.mdx b/docs/reference/configuration.mdx index da27cf8..05b87b6 100644 --- a/docs/reference/configuration.mdx +++ b/docs/reference/configuration.mdx @@ -60,6 +60,26 @@ Override the machine hostname sent with events. Useful in containerized environments where the system hostname is a random container ID. Set this to a stable identifier like the pod name or instance ID. +### `pg_stat_ch.extra_attributes` + +Exporter metadata expressed as semicolon-separated `key:value` pairs. `instance_uuid` is included in all export formats; the other resource columns are populated by the Arrow exporters. + +| | | +|---|---| +| **Type** | string | +| **Default** | `''` | +| **Context** | sighup | + +Supply `instance_uuid` with the Postgres service's canonical UUID alongside the existing `instance_ubid` and server attributes: + +```ini +pg_stat_ch.extra_attributes = 'instance_uuid:01234567-89ab-8ad0-9234-56789abcdef0;instance_ubid:pg04hmasw9ne4j8t5cy4tqkff1;server_role:primary' +``` + +The provisioner supplies the UUID; the producer does not derive it from the UBID. Native ClickHouse and Arrow exports include an `instance_uuid` string column, and ordinary OTLP exports include an `instance_uuid` log attribute. An omitted value produces an empty string. + +The bundled ClickHouse migration adds `instance_uuid String DEFAULT ''` to `events_raw`. Apply it **before deploying the producer when using native ClickHouse export**: its named inserts require the column even when the UUID is unset. ArrowStream inserts can arrive before the column exists; ClickHouse ignores the extra field until the schema includes it. External schemas and materialized-view projections need their own updates to retain the UUID; previously ignored values are not backfilled. + ### `pg_stat_ch.log_min_elevel` Minimum error severity level to capture via the `emit_log_hook`. diff --git a/docs/reference/events-schema.mdx b/docs/reference/events-schema.mdx index c163154..7859235 100644 --- a/docs/reference/events-schema.mdx +++ b/docs/reference/events-schema.mdx @@ -18,6 +18,7 @@ The table is partitioned by date (`toDate(ts_start)`) and ordered by `ts_start` | Column | Type | Description | |---|---|---| +| `instance_uuid` | `String` | Postgres service UUID supplied through `pg_stat_ch.extra_attributes`. Native ClickHouse and both Arrow exporters include this field; OTLP emits it as a log attribute. It is empty when unset. | | `db` | `LowCardinality(String)` | PostgreSQL database name. | | `username` | `LowCardinality(String)` | PostgreSQL user or role that executed the query. | | `pid` | `Int32` | Backend process ID. Correlate with `pg_stat_activity` for session-level debugging. | diff --git a/schema/migrations/20260925000001_add_instance_uuid.sql b/schema/migrations/20260925000001_add_instance_uuid.sql new file mode 100644 index 0000000..f388f5c --- /dev/null +++ b/schema/migrations/20260925000001_add_instance_uuid.sql @@ -0,0 +1,10 @@ +-- +goose Up +ALTER TABLE pg_stat_ch.events_raw + ADD COLUMN IF NOT EXISTS instance_uuid String DEFAULT '' + COMMENT 'Postgres service UUID supplied by pg_stat_ch.extra_attributes; empty when unset.' + AFTER instance_ubid; + +-- +goose Down +-- Retain the additive column on rollback: events_recent_1h's SELECT * depends +-- on it, and older producers can still insert without supplying it. +SELECT 1; diff --git a/src/config/guc.c b/src/config/guc.c index a9a8371..2728fe0 100644 --- a/src/config/guc.c +++ b/src/config/guc.c @@ -389,8 +389,9 @@ void PschInitGuc(void) { DefineCustomStringVariable( "pg_stat_ch.extra_attributes", - "Key-value pairs appended to exported Arrow batches.", - "Semicolon-separated k:v pairs for resource columns: " + "Key-value metadata for exported events.", + "Semicolon-separated k:v pairs for resource columns, including instance_uuid " + "for the Postgres service UUID in all export formats: " "'instance_ubid:abc;server_role:primary;read_replica_type:regional;region:us-east-1'.", &psch_extra_attributes, "", diff --git a/src/export/arrow_batch.cc b/src/export/arrow_batch.cc index 63bb650..2cafb5e 100644 --- a/src/export/arrow_batch.cc +++ b/src/export/arrow_batch.cc @@ -191,6 +191,7 @@ struct ArrowBatchBuilder::Impl { arrow::UInt32Builder parallel_workers_planned_builder; arrow::UInt32Builder parallel_workers_launched_builder; arrow::StringBuilder instance_ubid_builder; + arrow::StringBuilder instance_uuid_builder; arrow::StringBuilder server_ubid_builder; DictBuilder server_role_builder; DictBuilder read_replica_type_builder; @@ -259,6 +260,7 @@ struct ArrowBatchBuilder::Impl { arrow::field("parallel_workers_planned", arrow::uint32()), arrow::field("parallel_workers_launched", arrow::uint32()), arrow::field("instance_ubid", arrow::utf8()), + arrow::field("instance_uuid", arrow::utf8()), arrow::field("server_ubid", arrow::utf8()), arrow::field("server_role", DictionaryUtf8Type()), arrow::field("read_replica_type", DictionaryUtf8Type()), @@ -424,6 +426,8 @@ struct ArrowBatchBuilder::Impl { if (!AppendString(&instance_ubid_builder, ExtraAttr("instance_ubid"), "Arrow instance_ubid append") || + !AppendString(&instance_uuid_builder, ExtraAttr("instance_uuid"), + "Arrow instance_uuid append") || !AppendString(&server_ubid_builder, ExtraAttr("server_ubid"), "Arrow server_ubid append") || !AppendString(&server_role_builder, ExtraAttr("server_role"), "Arrow server_role append") || !AppendString(&read_replica_type_builder, ExtraAttr("read_replica_type"), @@ -436,13 +440,13 @@ struct ArrowBatchBuilder::Impl { return false; } - estimated_bytes += kFixedBytesPerRow + db_name.size() + db_user.size() + app.size() + - client_addr.size() + query_text.size() + err_message.size() + - err_sqlstate.size() + service_version.size() + - ExtraAttr("instance_ubid").size() + ExtraAttr("server_ubid").size() + - ExtraAttr("server_role").size() + ExtraAttr("read_replica_type").size() + - ExtraAttr("region").size() + ExtraAttr("cell").size() + - ExtraAttr("host_id").size() + ExtraAttr("pod_name").size(); + estimated_bytes += + kFixedBytesPerRow + db_name.size() + db_user.size() + app.size() + client_addr.size() + + query_text.size() + err_message.size() + err_sqlstate.size() + service_version.size() + + ExtraAttr("instance_ubid").size() + ExtraAttr("instance_uuid").size() + + ExtraAttr("server_ubid").size() + ExtraAttr("server_role").size() + + ExtraAttr("read_replica_type").size() + ExtraAttr("region").size() + + ExtraAttr("cell").size() + ExtraAttr("host_id").size() + ExtraAttr("pod_name").size(); ++num_rows; return true; } @@ -522,6 +526,7 @@ struct ArrowBatchBuilder::Impl { !add_array(¶llel_workers_planned_builder, "Arrow parallel_workers_planned finish") || !add_array(¶llel_workers_launched_builder, "Arrow parallel_workers_launched finish") || !add_array(&instance_ubid_builder, "Arrow instance_ubid finish") || + !add_array(&instance_uuid_builder, "Arrow instance_uuid finish") || !add_array(&server_ubid_builder, "Arrow server_ubid finish") || !add_dict_array(&server_role_builder, "Arrow server_role finish") || !add_dict_array(&read_replica_type_builder, "Arrow read_replica_type finish") || @@ -629,6 +634,7 @@ struct ArrowBatchBuilder::Impl { parallel_workers_planned_builder.Reset(); parallel_workers_launched_builder.Reset(); instance_ubid_builder.Reset(); + instance_uuid_builder.Reset(); server_ubid_builder.Reset(); server_role_builder.ResetFull(); read_replica_type_builder.ResetFull(); diff --git a/src/export/extra_attributes.h b/src/export/extra_attributes.h new file mode 100644 index 0000000..bc36595 --- /dev/null +++ b/src/export/extra_attributes.h @@ -0,0 +1,46 @@ +#ifndef PG_STAT_CH_SRC_EXPORT_EXTRA_ATTRIBUTES_H_ +#define PG_STAT_CH_SRC_EXPORT_EXTRA_ATTRIBUTES_H_ + +#include +#include +#include +#include + +// Parse "key1:val1;key2:val2" into a flat list. First match wins on +// duplicate keys (Get linear-scans from the front). Empty input -> empty list. +class ExtraAttrs { + public: + explicit ExtraAttrs(const char* raw) { + if (raw == nullptr) { + return; + } + std::string_view input(raw); + while (!input.empty()) { + const size_t delim = input.find(';'); + const std::string_view token = + (delim == std::string_view::npos) ? input : input.substr(0, delim); + const size_t sep = token.find(':'); + if (sep != std::string_view::npos) { + attrs_.emplace_back(std::string(token.substr(0, sep)), std::string(token.substr(sep + 1))); + } + if (delim == std::string_view::npos) { + break; + } + input.remove_prefix(delim + 1); + } + } + + std::string Get(std::string_view key) const { + for (const auto& [k, v] : attrs_) { + if (k == key) { + return v; + } + } + return {}; + } + + private: + std::vector> attrs_; +}; + +#endif // PG_STAT_CH_SRC_EXPORT_EXTRA_ATTRIBUTES_H_ diff --git a/src/export/otel_arrow_exporter.cc b/src/export/otel_arrow_exporter.cc index 369d876..5292f11 100644 --- a/src/export/otel_arrow_exporter.cc +++ b/src/export/otel_arrow_exporter.cc @@ -39,6 +39,7 @@ extern "C" { #include "pg_stat_ch/pg_stat_ch.h" #include "config/guc.h" #include "export/exporter_interface.h" +#include "export/extra_attributes.h" #include "export/otel_arrow_exporter.h" #include "export/otel_exporter.h" @@ -117,43 +118,6 @@ struct ArrowSlot { std::shared_ptr builder; }; -// Parse "key1:val1;key2:val2" into a flat list. First match wins on -// duplicate keys (Get linear-scans from the front). Empty input -> empty list. -class ExtraAttrs { - public: - explicit ExtraAttrs(const char* raw) { - if (raw == nullptr) { - return; - } - std::string_view input(raw); - while (!input.empty()) { - const size_t delim = input.find(';'); - const std::string_view token = - (delim == std::string_view::npos) ? input : input.substr(0, delim); - const size_t sep = token.find(':'); - if (sep != std::string_view::npos) { - attrs_.emplace_back(std::string(token.substr(0, sep)), std::string(token.substr(sep + 1))); - } - if (delim == std::string_view::npos) { - break; - } - input.remove_prefix(delim + 1); - } - } - - std::string Get(std::string_view key) const { - for (const auto& [k, v] : attrs_) { - if (k == key) { - return v; - } - } - return {}; - } - - private: - std::vector> attrs_; -}; - // --------------------------------------------------------------------------- class OTelArrowExporter : public StatsExporter { @@ -371,7 +335,7 @@ class OTelArrowExporter : public StatsExporter { // BeginRow so stats_exporter.cc's column-emission loop doesn't have to // know about them: // - // - 8 envelope columns + read_replica_type: per-process constants from + // - Envelope columns: per-process constants from // pg_stat_ch.extra_attributes (or "none" default for read_replica_type // per clickgres-platform's convention). // - service_version: PG_STAT_CH_VERSION macro, not from extra_attributes. diff --git a/src/export/stats_exporter.cc b/src/export/stats_exporter.cc index 6f38843..97198ce 100644 --- a/src/export/stats_exporter.cc +++ b/src/export/stats_exporter.cc @@ -18,6 +18,7 @@ extern "C" { #include "export/arrow_batch.h" #include "export/clickhouse_exporter.h" #include "export/exporter_interface.h" +#include "export/extra_attributes.h" #include "export/otel_arrow_exporter.h" #include "export/otel_exporter.h" #include "export/stats_exporter.h" @@ -239,6 +240,8 @@ void ExportEventStatsInternal(const std::vector& events, StatsExporte exporter->BeginBatch(); + const std::string instance_uuid = ExtraAttrs(psch_extra_attributes).Get("instance_uuid"); + auto col_instance_uuid = exporter->StatHCString("instance_uuid"); auto col_ts = exporter->StatTimestamp("ts"); auto col_duration_us = exporter->DbDurationColumn(); auto col_db_name = exporter->DbNameColumn(); @@ -294,6 +297,7 @@ void ExportEventStatsInternal(const std::vector& events, StatsExporte for (const auto& ev : events) { exporter->BeginRow(); + col_instance_uuid->Append(instance_uuid); col_ts->Append(ev.ts_start + kPostgresEpochOffsetUs); col_duration_us->Append(ev.duration_us); col_db_name->Append(std::string(ev.datname, ev.datname_len)); diff --git a/t/010_clickhouse_export.pl b/t/010_clickhouse_export.pl index 64a4bd4..4dfae40 100644 --- a/t/010_clickhouse_export.pl +++ b/t/010_clickhouse_export.pl @@ -63,6 +63,11 @@ 10 ); cmp_ok($query_check, '>=', 1, 'Query text is captured'); + + my $uuid_count = psch_query_clickhouse( + "SELECT count() FROM pg_stat_ch.events_raw WHERE instance_uuid != ''"); + chomp $uuid_count; + is($uuid_count, '0', 'Unconfigured instance UUID is empty'); }; # Test 2: Batch sizing - verify batch_max is honored @@ -112,6 +117,10 @@ # Test 4: All fields populated subtest 'all fields populated' => sub { + my $instance_uuid = '01234567-89ab-8ad0-9234-56789abcdef0'; + $node->safe_psql('postgres', + "ALTER SYSTEM SET pg_stat_ch.extra_attributes = 'instance_uuid:$instance_uuid'"); + $node->restart(); psch_query_clickhouse("TRUNCATE TABLE pg_stat_ch.events_raw"); psch_reset_stats($node); @@ -147,6 +156,13 @@ ); cmp_ok($db_operation_check, '>=', 1, 'db_operation is populated'); + my $uuid_check = psch_wait_for_clickhouse_query( + "SELECT count() FROM pg_stat_ch.events_raw WHERE instance_uuid = '$instance_uuid'", + sub { $_[0] >= 1 }, + 10 + ); + cmp_ok($uuid_check, '>=', 1, 'Configured instance UUID is exported natively'); + # Clean up $node->safe_psql('postgres', 'DROP TABLE IF EXISTS test_fields'); }; diff --git a/t/026_arrow_dump.pl b/t/026_arrow_dump.pl index ca0083d..0e0cdca 100644 --- a/t/026_arrow_dump.pl +++ b/t/026_arrow_dump.pl @@ -125,35 +125,39 @@ # Test 3: Validate IPC file contents with pyarrow (if available) # ============================================================================ SKIP: { - skip 'uv not installed (needed for pyarrow validation)', 1 unless $have_uv; - - subtest 'ipc file contents valid' => sub { - # Clean and produce a fresh dump. - unlink glob("$dump_dir/*.ipc"); - - psch_reset_stats($node); - - # Run a distinctive query we can look for. - $node->safe_psql('postgres', - 'CREATE TABLE IF NOT EXISTS arrow_test(id int)'); - $node->safe_psql('postgres', - "INSERT INTO arrow_test VALUES (42), (43), (44)"); - $node->safe_psql('postgres', 'SELECT * FROM arrow_test'); - $node->safe_psql('postgres', 'DROP TABLE arrow_test'); - - # Wait for dump. - my @ipc_files; - my $deadline = time() + 10; - while (time() < $deadline) { - @ipc_files = glob("$dump_dir/*.ipc"); - last if @ipc_files > 0; - select(undef, undef, undef, 0.2); - } - - cmp_ok(scalar @ipc_files, '>=', 1, 'IPC dump file present for validation'); - - # Validate with pyarrow via uv inline script. - my $validation_script = <<'PYEOF'; + skip 'uv not installed (needed for pyarrow validation)', 2 unless $have_uv; + + for my $instance_uuid ('', '01234567-89ab-8ad0-9234-56789abcdef0') { + subtest 'ipc file contents valid with ' . ($instance_uuid eq '' ? 'unset UUID' : 'configured UUID') => sub { + $node->safe_psql('postgres', + "ALTER SYSTEM SET pg_stat_ch.extra_attributes = 'instance_uuid:$instance_uuid'"); + $node->restart(); + # Clean and produce a fresh dump. + unlink glob("$dump_dir/*.ipc"); + + psch_reset_stats($node); + + # Run a distinctive query we can look for. + $node->safe_psql('postgres', + 'CREATE TABLE IF NOT EXISTS arrow_test(id int)'); + $node->safe_psql('postgres', + "INSERT INTO arrow_test VALUES (42), (43), (44)"); + $node->safe_psql('postgres', 'SELECT * FROM arrow_test'); + $node->safe_psql('postgres', 'DROP TABLE arrow_test'); + + # Wait for dump. + my @ipc_files; + my $deadline = time() + 10; + while (time() < $deadline) { + @ipc_files = glob("$dump_dir/*.ipc"); + last if @ipc_files > 0; + select(undef, undef, undef, 0.2); + } + + cmp_ok(scalar @ipc_files, '>=', 1, 'IPC dump file present for validation'); + + # Validate with pyarrow via uv inline script. + my $validation_script = <<'PYEOF'; # /// script # requires-python = ">=3.10" # dependencies = ["pyarrow"] @@ -166,13 +170,18 @@ total_rows = 0 schema = None -for path in sys.argv[1:]: +expected_uuid = sys.argv[1] + +for path in sys.argv[2:]: try: with open(path, 'rb') as f: reader = pa.ipc.open_stream(f) schema = reader.schema for batch in reader: total_rows += batch.num_rows + uuid_column = batch.column(schema.get_field_index("instance_uuid")) + assert uuid_column.type == pa.utf8(), f"instance_uuid type wrong: {uuid_column.type}" + assert all(value == expected_uuid for value in uuid_column.to_pylist()), "instance_uuid was not preserved" except Exception as e: errors.append(f"{path}: {e}") @@ -186,7 +195,7 @@ 'duration_us', 'rows', 'pid', 'query_id', 'shared_blks_hit', 'shared_blks_read', 'wal_records', 'wal_bytes', - 'service_version', 'region', 'read_replica_type', + 'service_version', 'region', 'read_replica_type', 'instance_uuid', ] missing = [c for c in expected if schema.get_field_index(c) == -1] if missing: @@ -211,27 +220,28 @@ print(f"OK:fields={len(schema)},rows={total_rows}") PYEOF - # Write script to temp file. - my $script_path = "$dump_dir/_validate.py"; - open(my $fh, '>', $script_path) or die "Cannot write $script_path: $!"; - print $fh $validation_script; - close $fh; - - my $file_args = join(' ', map { "'$_'" } @ipc_files); - my $raw_output = `uv run '$script_path' $file_args 2>&1`; - chomp($raw_output); - # uv prints install progress on earlier lines; grab the last line. - my @lines = split /\n/, $raw_output; - my $output = $lines[-1] // ''; - - like($output, qr/^OK:/, "pyarrow validation passed: $output"); - - # Extract row count and verify we captured events. - if ($output =~ /rows=(\d+)/) { - cmp_ok($1, '>=', 3, - "IPC files contain >= 3 rows (got $1)"); - } - }; + # Write script to temp file. + my $script_path = "$dump_dir/_validate.py"; + open(my $fh, '>', $script_path) or die "Cannot write $script_path: $!"; + print $fh $validation_script; + close $fh; + + my $file_args = join(' ', map { "'$_'" } @ipc_files); + my $raw_output = `uv run '$script_path' '$instance_uuid' $file_args 2>&1`; + chomp($raw_output); + # uv prints install progress on earlier lines; grab the last line. + my @lines = split /\n/, $raw_output; + my $output = $lines[-1] // ''; + + like($output, qr/^OK:/, "pyarrow validation passed: $output"); + + # Extract row count and verify we captured events. + if ($output =~ /rows=(\d+)/) { + cmp_ok($1, '>=', 3, + "IPC files contain >= 3 rows (got $1)"); + } + }; + } } # ============================================================================ diff --git a/t/036_unified_arrow_e2e.pl b/t/036_unified_arrow_e2e.pl index c0e9cf8..a780983 100644 --- a/t/036_unified_arrow_e2e.pl +++ b/t/036_unified_arrow_e2e.pl @@ -71,7 +71,7 @@ pg_stat_ch.use_unified_arrow_exporter = on pg_stat_ch.debug_arrow_dump_dir = '$dump_dir' pg_stat_ch.hostname = 'unified-arrow-e2e-host' -pg_stat_ch.extra_attributes = 'instance_ubid:test-instance;server_role:primary;region:test-region;cell:test-cell;read_replica_type:none' +pg_stat_ch.extra_attributes = 'instance_ubid:test-instance;instance_uuid:01234567-89ab-8ad0-9234-56789abcdef0;server_role:primary;region:test-region;cell:test-cell;read_replica_type:none' }); $node->start(); $node->safe_psql('postgres', 'CREATE EXTENSION pg_stat_ch'); @@ -180,10 +180,10 @@ # Envelope columns from extra_attributes were threaded through. my $envelope = psch_query_clickhouse( - "SELECT DISTINCT instance_ubid, server_role, region, cell, read_replica_type " . + "SELECT DISTINCT instance_ubid, instance_uuid, server_role, region, cell, read_replica_type " . "FROM pg_stat_ch.events_raw WHERE instance_ubid != '' LIMIT 1 FORMAT TSV"); chomp $envelope; -is($envelope, "test-instance\tprimary\ttest-region\ttest-cell\tnone", +is($envelope, "test-instance\t01234567-89ab-8ad0-9234-56789abcdef0\tprimary\ttest-region\ttest-cell\tnone", "envelope columns populated from extra_attributes"); $node->stop();