From 9e5cc2b28ce90e95e6c4c032771021945d60e453 Mon Sep 17 00:00:00 2001 From: amogiska Date: Fri, 25 Sep 2026 02:59:09 +0000 Subject: [PATCH 1/2] feat(arrow): include instance UUID in exported batches --- docs/reference/configuration.mdx | 18 +++ docs/reference/events-schema.mdx | 1 + .../20260925000001_add_instance_uuid.sql | 10 ++ src/config/guc.c | 3 +- src/export/arrow_batch.cc | 20 +-- src/export/otel_arrow_exporter.cc | 7 +- t/026_arrow_dump.pl | 114 ++++++++++-------- t/036_unified_arrow_e2e.pl | 6 +- 8 files changed, 115 insertions(+), 64 deletions(-) create mode 100644 schema/migrations/20260925000001_add_instance_uuid.sql diff --git a/docs/reference/configuration.mdx b/docs/reference/configuration.mdx index da27cf8..3acc231 100644 --- a/docs/reference/configuration.mdx +++ b/docs/reference/configuration.mdx @@ -60,6 +60,24 @@ 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` + +Resource columns included in both Arrow export formats, expressed as semicolon-separated `key:value` pairs. + +| | | +|---|---| +| **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 producer copies the UUID into a UTF-8 `instance_uuid` column; it does not derive it from the UBID. An omitted value produces an empty string. The provisioner must supply the UUID, and Arrow receivers must support the new column before the producer is rolled out. The bundled ClickHouse migration adds `instance_uuid String DEFAULT ''` to `events_raw`; external schemas need their own migration. + ### `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..f385a08 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`. Both Arrow exporters include this field; 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..67a97f2 100644 --- a/src/config/guc.c +++ b/src/config/guc.c @@ -390,7 +390,8 @@ 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: " + "Semicolon-separated k:v pairs for resource columns, including instance_uuid " + "for the Postgres service UUID: " "'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/otel_arrow_exporter.cc b/src/export/otel_arrow_exporter.cc index 369d876..51c0067 100644 --- a/src/export/otel_arrow_exporter.cc +++ b/src/export/otel_arrow_exporter.cc @@ -371,7 +371,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. @@ -403,6 +403,7 @@ class OTelArrowExporter : public StatsExporter { // Synthesized columns (populated implicitly in BeginRow). shared_ptr> inst_ubid_; + shared_ptr> inst_uuid_; shared_ptr> srv_ubid_; shared_ptr> srv_role_; shared_ptr> region_; @@ -414,6 +415,7 @@ class OTelArrowExporter : public StatsExporter { // Cached for per-row appends. std::string instance_ubid_val_; + std::string instance_uuid_val_; std::string server_ubid_val_; std::string server_role_val_; std::string region_val_; @@ -427,6 +429,7 @@ class OTelArrowExporter : public StatsExporter { void OTelArrowExporter::RegisterEnvelopeColumns() { // OTel resource attributes from psch_extra_attributes. inst_ubid_ = MakeUtf8Sv("instance_ubid"); + inst_uuid_ = MakeUtf8Sv("instance_uuid"); srv_ubid_ = MakeUtf8Sv("server_ubid"); srv_role_ = MakeDictSv("server_role"); read_replica_type_ = MakeDictSv("read_replica_type"); @@ -438,6 +441,7 @@ void OTelArrowExporter::RegisterEnvelopeColumns() { const ExtraAttrs attrs(psch_extra_attributes); instance_ubid_val_ = attrs.Get("instance_ubid"); + instance_uuid_val_ = attrs.Get("instance_uuid"); server_ubid_val_ = attrs.Get("server_ubid"); server_role_val_ = attrs.Get("server_role"); region_val_ = attrs.Get("region"); @@ -488,6 +492,7 @@ void OTelArrowExporter::BeginRow() { // Synthesized columns fire here so the call site doesn't need to know // about them. inst_ubid_->Append(instance_ubid_val_); + inst_uuid_->Append(instance_uuid_val_); srv_ubid_->Append(server_ubid_val_); srv_role_->Append(server_role_val_); read_replica_type_->Append(read_replica_type_val_); 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(); From 6a9aafdc0d9eff31b8e495eb375d0a24b6c2d127 Mon Sep 17 00:00:00 2001 From: amogiska Date: Fri, 25 Sep 2026 07:03:54 +0000 Subject: [PATCH 2/2] fix(exporter): emit instance UUID through shared columns --- docs/reference/configuration.mdx | 6 ++-- docs/reference/events-schema.mdx | 2 +- src/config/guc.c | 4 +-- src/export/extra_attributes.h | 46 +++++++++++++++++++++++++++++++ src/export/otel_arrow_exporter.cc | 43 +---------------------------- src/export/stats_exporter.cc | 4 +++ t/010_clickhouse_export.pl | 16 +++++++++++ 7 files changed, 74 insertions(+), 47 deletions(-) create mode 100644 src/export/extra_attributes.h diff --git a/docs/reference/configuration.mdx b/docs/reference/configuration.mdx index 3acc231..05b87b6 100644 --- a/docs/reference/configuration.mdx +++ b/docs/reference/configuration.mdx @@ -62,7 +62,7 @@ Useful in containerized environments where the system hostname is a random conta ### `pg_stat_ch.extra_attributes` -Resource columns included in both Arrow export formats, expressed as semicolon-separated `key:value` pairs. +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. | | | |---|---| @@ -76,7 +76,9 @@ Supply `instance_uuid` with the Postgres service's canonical UUID alongside the pg_stat_ch.extra_attributes = 'instance_uuid:01234567-89ab-8ad0-9234-56789abcdef0;instance_ubid:pg04hmasw9ne4j8t5cy4tqkff1;server_role:primary' ``` -The producer copies the UUID into a UTF-8 `instance_uuid` column; it does not derive it from the UBID. An omitted value produces an empty string. The provisioner must supply the UUID, and Arrow receivers must support the new column before the producer is rolled out. The bundled ClickHouse migration adds `instance_uuid String DEFAULT ''` to `events_raw`; external schemas need their own migration. +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` diff --git a/docs/reference/events-schema.mdx b/docs/reference/events-schema.mdx index f385a08..7859235 100644 --- a/docs/reference/events-schema.mdx +++ b/docs/reference/events-schema.mdx @@ -18,7 +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`. Both Arrow exporters include this field; it is empty when unset. | +| `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/src/config/guc.c b/src/config/guc.c index 67a97f2..2728fe0 100644 --- a/src/config/guc.c +++ b/src/config/guc.c @@ -389,9 +389,9 @@ void PschInitGuc(void) { DefineCustomStringVariable( "pg_stat_ch.extra_attributes", - "Key-value pairs appended to exported Arrow batches.", + "Key-value metadata for exported events.", "Semicolon-separated k:v pairs for resource columns, including instance_uuid " - "for the Postgres service 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/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 51c0067..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 { @@ -403,7 +367,6 @@ class OTelArrowExporter : public StatsExporter { // Synthesized columns (populated implicitly in BeginRow). shared_ptr> inst_ubid_; - shared_ptr> inst_uuid_; shared_ptr> srv_ubid_; shared_ptr> srv_role_; shared_ptr> region_; @@ -415,7 +378,6 @@ class OTelArrowExporter : public StatsExporter { // Cached for per-row appends. std::string instance_ubid_val_; - std::string instance_uuid_val_; std::string server_ubid_val_; std::string server_role_val_; std::string region_val_; @@ -429,7 +391,6 @@ class OTelArrowExporter : public StatsExporter { void OTelArrowExporter::RegisterEnvelopeColumns() { // OTel resource attributes from psch_extra_attributes. inst_ubid_ = MakeUtf8Sv("instance_ubid"); - inst_uuid_ = MakeUtf8Sv("instance_uuid"); srv_ubid_ = MakeUtf8Sv("server_ubid"); srv_role_ = MakeDictSv("server_role"); read_replica_type_ = MakeDictSv("read_replica_type"); @@ -441,7 +402,6 @@ void OTelArrowExporter::RegisterEnvelopeColumns() { const ExtraAttrs attrs(psch_extra_attributes); instance_ubid_val_ = attrs.Get("instance_ubid"); - instance_uuid_val_ = attrs.Get("instance_uuid"); server_ubid_val_ = attrs.Get("server_ubid"); server_role_val_ = attrs.Get("server_role"); region_val_ = attrs.Get("region"); @@ -492,7 +452,6 @@ void OTelArrowExporter::BeginRow() { // Synthesized columns fire here so the call site doesn't need to know // about them. inst_ubid_->Append(instance_ubid_val_); - inst_uuid_->Append(instance_uuid_val_); srv_ubid_->Append(server_ubid_val_); srv_role_->Append(server_role_val_); read_replica_type_->Append(read_replica_type_val_); 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'); };