From ec5184a4535c13b02e3dbf4bf1f533eba54054a3 Mon Sep 17 00:00:00 2001 From: Lalit Maganti Date: Thu, 1 Oct 2026 05:39:08 +0100 Subject: [PATCH 1/3] tp: make FlexVector::resize grow geometrically FlexVector::resize allocated exactly the size asked for, so growing a buffer a batch at a time with resize() reallocated and copied it every time: quadratic. BitVector::resize, built on it, did the same. Code worked around it with hand-rolled doubling, or by calling reserve(), which already grew geometrically, before every resize(). resize() now grows capacity as reserve() does: to at least 1.5x when it has to reallocate. A resize of an empty vector still allocates exactly what is asked for. The workarounds in the tree operators go, as does BitVector::reserve(), which existed only for them. --- .../core/exec/tree_number_nodes.cc | 4 +--- src/trace_processor/core/exec/tree_order.cc | 16 ++++++---------- src/trace_processor/core/util/bit_vector.h | 5 ----- src/trace_processor/core/util/flex_vector.h | 11 ++--------- .../core/util/flex_vector_unittest.cc | 8 ++++++++ 5 files changed, 17 insertions(+), 27 deletions(-) diff --git a/src/trace_processor/core/exec/tree_number_nodes.cc b/src/trace_processor/core/exec/tree_number_nodes.cc index 2b5b6e70b59..e5ed7d15f57 100644 --- a/src/trace_processor/core/exec/tree_number_nodes.cc +++ b/src/trace_processor/core/exec/tree_number_nodes.cc @@ -16,7 +16,6 @@ #include "src/trace_processor/core/exec/tree_number_nodes.h" -#include #include #include @@ -268,8 +267,7 @@ bool TreeNumberNodes::NumberByKey(const RowBatch& in, return false; } if (s.has_row.size() <= node) { - // Grown geometrically, as resize allocates exactly what it is asked for. - s.has_row.resize(std::max(node + 1, s.has_row.size() * 2)); + s.has_row.resize(node + 1); } if (s.has_row.is_set(node)) { s.status = base::ErrStatus( diff --git a/src/trace_processor/core/exec/tree_order.cc b/src/trace_processor/core/exec/tree_order.cc index cfa5c958208..61d65ef4fcb 100644 --- a/src/trace_processor/core/exec/tree_order.cc +++ b/src/trace_processor/core/exec/tree_order.cc @@ -100,10 +100,8 @@ bool TreeChildFirst::Consume(const RowBatch& in, Breaker::State& state) const { s.nodes_seen = std::max(s.nodes_seen, parent + 1); } if (s.has_row.size() < s.nodes_seen) { - // Grown geometrically, as resize allocates exactly what it is asked for. - auto size = std::max(s.nodes_seen, s.has_row.size() * 2); - s.has_row.resize(size); - s.row_of_node.resize(size); + s.has_row.resize(s.nodes_seen); + s.row_of_node.resize(s.nodes_seen); } if (s.has_row.is_set(node)) { s.status = base::ErrStatus("%s: more than one row has the same node", @@ -253,12 +251,10 @@ void TreeParentFirst::State::Nodes::Grow(uint32_t new_count) { if (new_count <= size) { return; } - // Grown geometrically, as resize allocates exactly what it is asked for. - uint32_t new_size = std::max(new_count, size * 2); - has_row.resize(new_size); - out.resize(new_size); - first_waiting.resize(new_size); - std::fill(first_waiting.data() + size, first_waiting.data() + new_size, + has_row.resize(new_count); + out.resize(new_count); + first_waiting.resize(new_count); + std::fill(first_waiting.data() + size, first_waiting.data() + new_count, kNoNode); } diff --git a/src/trace_processor/core/util/bit_vector.h b/src/trace_processor/core/util/bit_vector.h index 42103df1a3c..56cd3ca0e1e 100644 --- a/src/trace_processor/core/util/bit_vector.h +++ b/src/trace_processor/core/util/bit_vector.h @@ -321,11 +321,6 @@ struct BitVector { memset(words_.data(), 0, words_.size() * sizeof(uint64_t)); } - // Makes room for `new_size` bits without changing the size. Growth here is - // geometric where resize allocates exactly what was asked for, so anything - // filling a bit vector a chunk at a time has to come through here first. - void reserve(uint64_t new_size) { words_.reserve((new_size + 63) / 64); } - // Resizes the vector to the specified size. If shrinking, bits past the new // size are cleared. If growing, new bits are set to the given value. void resize(uint64_t new_size, bool value = false) { diff --git a/src/trace_processor/core/util/flex_vector.h b/src/trace_processor/core/util/flex_vector.h index 4014fe18808..f2e108b973b 100644 --- a/src/trace_processor/core/util/flex_vector.h +++ b/src/trace_processor/core/util/flex_vector.h @@ -172,16 +172,9 @@ class FlexVector { } // Resizes the vector to the specified size. If growing, new elements are - // uninitialized. + // uninitialized. Capacity grows as with reserve(). void resize(uint64_t new_size) { - if (new_size > capacity()) { - Slab new_slab = - Slab::Alloc(base::AlignUp(new_size, kCapacityMultiple)); - if (size_ > 0) { - memcpy(new_slab.data(), slab_.data(), size_ * sizeof(T)); - } - slab_ = std::move(new_slab); - } + reserve(new_size); size_ = new_size; } diff --git a/src/trace_processor/core/util/flex_vector_unittest.cc b/src/trace_processor/core/util/flex_vector_unittest.cc index 19afaff11bf..4f1da2df2af 100644 --- a/src/trace_processor/core/util/flex_vector_unittest.cc +++ b/src/trace_processor/core/util/flex_vector_unittest.cc @@ -265,5 +265,13 @@ TEST(FlexVectorTest, ResizeGrow) { EXPECT_EQ(vec[2], 3); } +TEST(FlexVectorTest, ResizeGrowsGeometrically) { + FlexVector vec; + vec.resize(1000); + uint64_t capacity = vec.capacity(); + vec.resize(capacity + 1); + EXPECT_GE(vec.capacity(), capacity * 3 / 2); +} + } // namespace } // namespace perfetto::trace_processor::core From 80ada084c5661759643ce49370308f0bb520a69a Mon Sep 17 00:00:00 2001 From: Lalit Maganti Date: Wed, 30 Sep 2026 00:52:14 +0100 Subject: [PATCH 2/3] tp: let INTERVAL INTERSECTION's PER keys be of any one type PER keys had to be integers and were widened to Int64 on every batch. They are now laid out as-is by a KeyEncoder, built on RowLayout, and compared as bytes. Integers of every width, doubles and strings can all be keys; integers of different widths agree, and nulls agree as under GROUP BY. A key must hold the same type in every operand. KeyEncoder works out each key column's type on the first batch and picks a writer for it; later batches only check that the columns are shaped as before. RowLayout::Slot is 8 bytes and passed by value, and rows are written through a PERFETTO_RESTRICT pointer, so the slot stays in registers while a column is written (2048 Int64 keys: 970 -> 708 ns; one key: 6.1 -> 4.3 ns). FlatColumnReader, beside ColumnView, reads a flat column through its selection and validity, for INTERVAL INTERSECTION's ts and dur and for string keys. --- Android.bp | 2 + BUILD | 2 + include/perfetto/base/compiler.h | 11 + src/trace_processor/core/common/row_layout.h | 18 +- src/trace_processor/core/exec/BUILD.gn | 5 + src/trace_processor/core/exec/column_view.h | 25 ++ .../core/exec/interval_intersect.cc | 85 +++---- .../core/exec/interval_intersect.h | 2 +- src/trace_processor/core/exec/key_encoder.cc | 216 ++++++++++++++++++ src/trace_processor/core/exec/key_encoder.h | 79 +++++++ .../core/exec/key_encoder_benchmark.cc | 87 +++++++ .../core/exec/key_encoder_unittest.cc | 99 ++++++++ .../perfetto_sql_connection_unittest.cc | 27 +++ .../perfetto_sql/pipeline/physical_plan.cc | 7 +- 14 files changed, 594 insertions(+), 71 deletions(-) create mode 100644 src/trace_processor/core/exec/key_encoder.cc create mode 100644 src/trace_processor/core/exec/key_encoder.h create mode 100644 src/trace_processor/core/exec/key_encoder_benchmark.cc create mode 100644 src/trace_processor/core/exec/key_encoder_unittest.cc diff --git a/Android.bp b/Android.bp index a1c4ede2dcc..48456d4146e 100644 --- a/Android.bp +++ b/Android.bp @@ -17823,6 +17823,7 @@ filegroup { "src/trace_processor/core/exec/column_view.cc", "src/trace_processor/core/exec/dataframe_scan.cc", "src/trace_processor/core/exec/interval_intersect.cc", + "src/trace_processor/core/exec/key_encoder.cc", "src/trace_processor/core/exec/operator.cc", "src/trace_processor/core/exec/pipeline.cc", "src/trace_processor/core/exec/row_batch.cc", @@ -17848,6 +17849,7 @@ filegroup { "src/trace_processor/core/exec/dataframe_scan_unittest.cc", "src/trace_processor/core/exec/executor_contract_unittest.cc", "src/trace_processor/core/exec/interval_intersect_unittest.cc", + "src/trace_processor/core/exec/key_encoder_unittest.cc", "src/trace_processor/core/exec/operator_unittest.cc", "src/trace_processor/core/exec/row_batch_unittest.cc", "src/trace_processor/core/exec/row_store_unittest.cc", diff --git a/BUILD b/BUILD index cbbcb8f39ea..cab143f6857 100644 --- a/BUILD +++ b/BUILD @@ -2587,6 +2587,8 @@ perfetto_filegroup( "src/trace_processor/core/exec/dataframe_scan.h", "src/trace_processor/core/exec/interval_intersect.cc", "src/trace_processor/core/exec/interval_intersect.h", + "src/trace_processor/core/exec/key_encoder.cc", + "src/trace_processor/core/exec/key_encoder.h", "src/trace_processor/core/exec/operator.cc", "src/trace_processor/core/exec/operator.h", "src/trace_processor/core/exec/pipeline.cc", diff --git a/include/perfetto/base/compiler.h b/include/perfetto/base/compiler.h index 8ffa94a0507..d07e733a6be 100644 --- a/include/perfetto/base/compiler.h +++ b/include/perfetto/base/compiler.h @@ -61,6 +61,17 @@ #define PERFETTO_NORETURN __declspec(noreturn) #endif +// Promises that what a pointer reaches is reached through nothing else in its +// scope, so writes through it cannot change what other pointers read. Not +// `__restrict`: the macOS SDK defines that away to nothing in C++. +#if defined(__GNUC__) || defined(__clang__) +#define PERFETTO_RESTRICT __restrict__ +#elif defined(_MSC_VER) +#define PERFETTO_RESTRICT __restrict +#else +#define PERFETTO_RESTRICT +#endif + #if defined(__GNUC__) || defined(__clang__) #define PERFETTO_DEBUG_FUNCTION_IDENTIFIER() __PRETTY_FUNCTION__ #elif defined(_MSC_VER) diff --git a/src/trace_processor/core/common/row_layout.h b/src/trace_processor/core/common/row_layout.h index 85766ba7a9f..59c77c39766 100644 --- a/src/trace_processor/core/common/row_layout.h +++ b/src/trace_processor/core/common/row_layout.h @@ -19,6 +19,7 @@ #include #include +#include #include #include @@ -39,21 +40,26 @@ class RowLayout { bool nullable = false; bool descending = false; }; + // Small enough to pass in one register: Write() takes it by value, so its + // fields stay in registers while row bytes are written. struct Slot { - uint32_t offset; - uint32_t stride; + uint16_t offset; + uint16_t stride; bool nullable; bool descending; }; + static_assert(sizeof(Slot) <= 8); RowLayout() = default; explicit RowLayout(const std::vector& columns) { for (const Column& column : columns) { - slots_.push_back({stride_, 0, column.nullable, column.descending}); + slots_.push_back({static_cast(stride_), 0, column.nullable, + column.descending}); stride_ += (column.nullable ? 1u : 0u) + ValueSize(column.type); } + PERFETTO_CHECK(stride_ <= std::numeric_limits::max()); for (Slot& slot : slots_) { - slot.stride = stride_; + slot.stride = static_cast(stride_); } } @@ -63,10 +69,10 @@ class RowLayout { // `get(i, &value)` returns false if row i has no value. template - PERFETTO_ALWAYS_INLINE static void Write(const Slot& slot, + PERFETTO_ALWAYS_INLINE static void Write(Slot slot, uint32_t count, Get get, - uint8_t* rows) { + uint8_t* PERFETTO_RESTRICT rows) { // Any other type would be implicitly converted to one of these by // Encode, writing a value of a different size to its slot. static_assert(std::is_same_v || std::is_same_v || diff --git a/src/trace_processor/core/exec/BUILD.gn b/src/trace_processor/core/exec/BUILD.gn index cecbb41b05c..f0e9c53733c 100644 --- a/src/trace_processor/core/exec/BUILD.gn +++ b/src/trace_processor/core/exec/BUILD.gn @@ -32,6 +32,8 @@ source_set("exec") { "dataframe_scan.h", "interval_intersect.cc", "interval_intersect.h", + "key_encoder.cc", + "key_encoder.h", "operator.cc", "operator.h", "pipeline.cc", @@ -80,6 +82,7 @@ perfetto_unittest_source_set("unittests") { "dataframe_scan_unittest.cc", "executor_contract_unittest.cc", "interval_intersect_unittest.cc", + "key_encoder_unittest.cc", "operator_unittest.cc", "row_batch_unittest.cc", "row_store_unittest.cc", @@ -105,6 +108,7 @@ if (enable_perfetto_benchmarks) { source_set("benchmarks") { testonly = true sources = [ + "key_encoder_benchmark.cc", "tree_number_nodes_benchmark.cc", "tree_order_benchmark.cc", ] @@ -115,6 +119,7 @@ if (enable_perfetto_benchmarks) { "../../containers", "../common", "../dataframe", + "../util", ] } } diff --git a/src/trace_processor/core/exec/column_view.h b/src/trace_processor/core/exec/column_view.h index 344cb7a499b..f69130437e3 100644 --- a/src/trace_processor/core/exec/column_view.h +++ b/src/trace_processor/core/exec/column_view.h @@ -131,6 +131,31 @@ class ColumnView { const BitVector* validity_ = nullptr; }; +// Reads a flat column's values through its selection and validity. +template +class FlatColumnReader { + public: + explicit FlatColumnReader(const ColumnView& column) + : data_(static_cast(column.data())), + selection_(column.selection()), + validity_(column.validity()) {} + + // False if the row holds no value. + PERFETTO_ALWAYS_INLINE bool Read(uint32_t row, T* out) const { + uint32_t index = selection_.GetIndex(row); + if (validity_ && !validity_->is_set(index)) { + return false; + } + *out = data_[index]; + return true; + } + + private: + const T* data_; + RowSelection selection_; + const BitVector* validity_; +}; + // Whether two batches' views of a column can be combined. An implicit Id and // a stored Uint32 hold the same values: gathering turns one into the other. inline bool SameLogicalType(const ColumnView& a, const ColumnView& b) { diff --git a/src/trace_processor/core/exec/interval_intersect.cc b/src/trace_processor/core/exec/interval_intersect.cc index 7ac67e27f5d..4d6afa7e6b8 100644 --- a/src/trace_processor/core/exec/interval_intersect.cc +++ b/src/trace_processor/core/exec/interval_intersect.cc @@ -21,7 +21,9 @@ #include #include #include +#include #include +#include #include #include @@ -35,6 +37,7 @@ #include "src/trace_processor/core/common/storage_types.h" #include "src/trace_processor/core/exec/column_chunk.h" #include "src/trace_processor/core/exec/column_view.h" +#include "src/trace_processor/core/exec/key_encoder.h" #include "src/trace_processor/core/exec/operator.h" #include "src/trace_processor/core/exec/row_batch.h" #include "src/trace_processor/core/exec/row_selection.h" @@ -85,27 +88,6 @@ struct Group { bool nonoverlapping = true; }; -// Reads one flat, non-null Int64 column of a batch. -struct Int64Reader { - explicit Int64Reader(const ColumnView& column) - : data(static_cast(column.data())), - selection(column.selection()), - validity(column.validity()) {} - - bool Read(uint32_t row, int64_t* out) const { - uint32_t index = selection.GetIndex(row); - if (validity && !validity->is_set(index)) { - return false; - } - *out = data[index]; - return true; - } - - const int64_t* data; - RowSelection selection; - const BitVector* validity; -}; - base::Status ValidateOperand(const RowBatch& batch, const IntervalIntersectOperand& operand, uint32_t which) { @@ -114,39 +96,25 @@ base::Status ValidateOperand(const RowBatch& batch, return column.kind() == ColumnView::Kind::kFlat && column.type().Is(); }; - bool ok = is_int64(operand.ts_column) && is_int64(operand.dur_column) && - std::all_of(operand.key_columns.begin(), operand.key_columns.end(), - is_int64); + bool ok = is_int64(operand.ts_column) && is_int64(operand.dur_column); return ok ? base::OkStatus() : base::ErrStatus( - "INTERVAL INTERSECTION: operand %u's ts, dur and PER " - "columns must be Int64", + "INTERVAL INTERSECTION: operand %u's ts and dur columns " + "must be Int64", which + 1); } -// A key column takes this many bytes: whether it holds a value, then the -// value itself. Two rows which hold no value there agree on it, as they -// would under GROUP BY. -constexpr size_t kKeyColumnBytes = 1 + sizeof(int64_t); - -void WriteKeyColumn(std::string& key, - uint32_t at, - bool present, - int64_t value) { - char* to = key.data() + at * kKeyColumnBytes; - to[0] = present ? 1 : 0; - memcpy(to + 1, &value, sizeof(value)); -} - class IntersectState : public OperatorState { public: ~IntersectState() override; std::vector> operand_states; std::vector> stores; + // Shared by the operands, so their keys' types must agree. + KeyEncoder keys; // Per operand, its rows by key. A key with no entry in some operand covers // nothing, so only keys every operand has produce regions. - std::vector> groups; + std::vector> groups; bool computed = false; Regions regions; @@ -181,17 +149,19 @@ base::Status Collect(const IntervalIntersectOperand& operand, uint32_t which, OperatorState& state, RowStore& store, - base::FlatHashMap& groups) { + KeyEncoder& keys, + base::FlatHashMapV2& groups) { RowBatch batch; RowBatch retained; - std::string key(operand.key_columns.size() * kKeyColumnBytes, '\0'); while (operand.source->GetData(batch, state)) { RETURN_IF_ERROR(ValidateOperand(batch, operand, which)); - Int64Reader ts(batch.column(operand.ts_column)); - Int64Reader dur(batch.column(operand.dur_column)); - std::vector keys; - for (uint32_t column : operand.key_columns) { - keys.emplace_back(batch.column(column)); + FlatColumnReader ts(batch.column(operand.ts_column)); + FlatColumnReader dur(batch.column(operand.dur_column)); + if (std::optional bad = keys.Encode(batch, operand.key_columns)) { + return base::ErrStatus( + "INTERVAL INTERSECTION: operand %u's PER column %u must hold one " + "type, the same in every operand", + which + 1, *bad + 1); } for (uint32_t row = 0; row < batch.size(); ++row) { int64_t start; @@ -205,17 +175,14 @@ base::Status Collect(const IntervalIntersectOperand& operand, "below zero", which + 1); } - for (uint32_t i = 0; i < keys.size(); ++i) { - // Read before writing: argument evaluation order is unspecified, so - // passing both `Read(&value)` and `value` could copy it unread. - int64_t value = 0; - bool present = keys[i].Read(row, &value); - WriteKeyColumn(key, i, present, value); + std::string_view key = keys.Key(row); + Group* group = groups.Find(key); + if (!group) { + group = groups.Insert(std::string(key), Group{}).first; } - Group& group = groups[key]; - group.intervals.push_back({static_cast(start), - static_cast(start + length), - store.size() + row}); + group->intervals.push_back({static_cast(start), + static_cast(start + length), + store.size() + row}); } retained.Reset(); for (uint32_t column : operand.retained_columns) { @@ -356,7 +323,7 @@ bool IntervalIntersect::GetData(RowBatch& out, OperatorState& state) const { s.scratch.narrowed.operands = count; for (uint32_t i = 0; i < count; ++i) { s.status = Collect(operands_[i], i, *s.operand_states[i], *s.stores[i], - s.groups[i]); + s.keys, s.groups[i]); if (!s.status.ok()) { return false; } diff --git a/src/trace_processor/core/exec/interval_intersect.h b/src/trace_processor/core/exec/interval_intersect.h index a233bcfb1e5..dda7e9fe513 100644 --- a/src/trace_processor/core/exec/interval_intersect.h +++ b/src/trace_processor/core/exec/interval_intersect.h @@ -32,7 +32,7 @@ namespace perfetto::trace_processor::core::exec { struct IntervalIntersectOperand { // Read in full before any region is found. const Source* source = nullptr; - // All must be flat Int64 columns of `source`'s batches. + // Flat Int64 columns of `source`'s batches. uint32_t ts_column = 0; uint32_t dur_column = 0; // Compared pairwise across the operands, in this order. diff --git a/src/trace_processor/core/exec/key_encoder.cc b/src/trace_processor/core/exec/key_encoder.cc new file mode 100644 index 00000000000..f1b2faec079 --- /dev/null +++ b/src/trace_processor/core/exec/key_encoder.cc @@ -0,0 +1,216 @@ +/* + * Copyright (C) 2026 The Android Open Source Project + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +#include "src/trace_processor/core/exec/key_encoder.h" + +#include +#include +#include +#include + +#include "perfetto/base/compiler.h" +#include "perfetto/base/logging.h" +#include "src/trace_processor/containers/string_pool.h" +#include "src/trace_processor/core/common/row_layout.h" +#include "src/trace_processor/core/common/storage_types.h" +#include "src/trace_processor/core/exec/column_view.h" +#include "src/trace_processor/core/exec/row_batch.h" +#include "src/trace_processor/core/exec/row_selection.h" +#include "src/trace_processor/core/util/bit_vector.h" + +namespace perfetto::trace_processor::core::exec { +namespace { + +template +PERFETTO_ALWAYS_INLINE void Write(const ColumnView& column, + RowLayout::Slot slot, + uint32_t count, + uint8_t* rows, + Value value) { + RowSelection selection = column.selection(); + const BitVector* validity = column.validity(); + RowLayout::Write( + slot, count, + [&](uint32_t row, T* out) { + uint32_t index = selection.GetIndex(row); + if (validity && !validity->is_set(index)) { + return false; + } + *out = value(index); + return true; + }, + rows); +} + +void WriteSequence(const ColumnView& column, + RowLayout::Slot slot, + uint32_t count, + uint8_t* rows) { + Write(column, slot, count, rows, + [](uint32_t index) { return int64_t{index}; }); +} + +template +void WriteInteger(const ColumnView& column, + RowLayout::Slot slot, + uint32_t count, + uint8_t* rows) { + const auto* data = static_cast(column.data()); + Write(column, slot, count, rows, + [data](uint32_t index) { return int64_t{data[index]}; }); +} + +void WriteDouble(const ColumnView& column, + RowLayout::Slot slot, + uint32_t count, + uint8_t* rows) { + const auto* data = static_cast(column.data()); + Write(column, slot, count, rows, + [data](uint32_t index) { return data[index]; }); +} + +// Not nullable: the null id is 0. +void WriteString(const ColumnView& column, + RowLayout::Slot slot, + uint32_t count, + uint8_t* rows) { + FlatColumnReader ids(column); + RowLayout::Write( + slot, count, + [&](uint32_t row, uint32_t* out) { + StringPool::Id id; + *out = ids.Read(row, &id) ? id.raw_id() : 0; + return true; + }, + rows); +} + +uint16_t ShapeOf(const ColumnView& column) { + return static_cast(static_cast(column.kind()) << 8 | + column.type().index()); +} + +} // namespace + +std::optional KeyEncoder::Encode( + const RowBatch& batch, + const std::vector& columns) { + bool same = key_columns_.size() == columns.size(); + for (uint32_t k = 0; same && k < columns.size(); ++k) { + same = key_columns_[k].shape == ShapeOf(batch.column(columns[k])); + } + if (PERFETTO_UNLIKELY(!same)) { + return Classify(batch, columns); + } + WriteKeys(batch, columns); + return std::nullopt; +} + +PERFETTO_ALWAYS_INLINE void KeyEncoder::WriteKeys( + const RowBatch& batch, + const std::vector& columns) { + uint32_t count = batch.size(); + bytes_.resize(size_t{layout_.stride()} * count); + uint8_t* rows = bytes_.data(); + const uint32_t* index = columns.data(); + for (const KeyColumn *c = key_columns_.data(), *end = c + key_columns_.size(); + c != end; ++c, ++index) { + c->writer(batch.column(*index), c->slot, count, rows); + } +} + +PERFETTO_NO_INLINE std::optional KeyEncoder::Classify( + const RowBatch& batch, + const std::vector& columns) { + bool first = kinds_.empty(); + kinds_.resize(columns.size()); + key_columns_.resize(columns.size()); + for (uint32_t k = 0; k < columns.size(); ++k) { + const ColumnView& column = batch.column(columns[k]); + std::optional kind; + Writer writer = nullptr; + switch (column.kind()) { + case ColumnView::Kind::kSequence: + kind = Kind::kInteger; + writer = &WriteSequence; + break; + case ColumnView::Kind::kFlat: + switch (column.type().index()) { + case StorageType::GetTypeIndex(): + kind = Kind::kInteger; + writer = &WriteInteger; + break; + case StorageType::GetTypeIndex(): + kind = Kind::kInteger; + writer = &WriteInteger; + break; + case StorageType::GetTypeIndex(): + kind = Kind::kInteger; + writer = &WriteInteger; + break; + case StorageType::GetTypeIndex(): + kind = Kind::kDouble; + writer = &WriteDouble; + break; + case StorageType::GetTypeIndex(): + kind = Kind::kString; + writer = &WriteString; + break; + default: + break; + } + break; + case ColumnView::Kind::kVariant: + break; + } + if (!kind || (!first && kinds_[k] != *kind)) { + // Nothing is kept of a batch which fails: the next is checked in full, + // and laid out if none has been yet. + key_columns_.clear(); + if (first) { + kinds_.clear(); + } + return k; + } + kinds_[k] = *kind; + key_columns_[k].shape = ShapeOf(column); + key_columns_[k].writer = writer; + } + if (first) { + std::vector layout; + for (Kind kind : kinds_) { + switch (kind) { + case Kind::kInteger: + layout.push_back({RowLayout::Type::kInt64, true}); + break; + case Kind::kDouble: + layout.push_back({RowLayout::Type::kDouble, true}); + break; + case Kind::kString: + layout.push_back({RowLayout::Type::kUint32, false}); + break; + } + } + layout_ = RowLayout(layout); + } + for (uint32_t k = 0; k < key_columns_.size(); ++k) { + key_columns_[k].slot = layout_.slot(k); + } + WriteKeys(batch, columns); + return std::nullopt; +} + +} // namespace perfetto::trace_processor::core::exec diff --git a/src/trace_processor/core/exec/key_encoder.h b/src/trace_processor/core/exec/key_encoder.h new file mode 100644 index 00000000000..a47f3a962cd --- /dev/null +++ b/src/trace_processor/core/exec/key_encoder.h @@ -0,0 +1,79 @@ +/* + * Copyright (C) 2026 The Android Open Source Project + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +#ifndef SRC_TRACE_PROCESSOR_CORE_EXEC_KEY_ENCODER_H_ +#define SRC_TRACE_PROCESSOR_CORE_EXEC_KEY_ENCODER_H_ + +#include +#include +#include +#include +#include + +#include "src/trace_processor/core/common/row_layout.h" +#include "src/trace_processor/core/exec/row_batch.h" +#include "src/trace_processor/core/util/flex_vector.h" + +namespace perfetto::trace_processor::core::exec { + +// Lays out rows' keys so rows with equal keys have equal bytes. Integers of +// every width are laid out as Int64 so they agree, and nulls agree with each +// other as under GROUP BY. A column must hold one type in every batch. +class KeyEncoder { + public: + // Returns the position in `columns` of one which cannot be a key or whose + // type changed. + std::optional Encode(const RowBatch& batch, + const std::vector& columns); + + std::string_view Key(uint32_t row) const { + return {reinterpret_cast(bytes_.data()) + + size_t{row} * layout_.stride(), + layout_.stride()}; + } + + private: + enum class Kind : uint8_t { kInteger, kDouble, kString }; + // Lays out one column's keys, for columns of one storage type. + using Writer = void (*)(const ColumnView& column, + RowLayout::Slot slot, + uint32_t count, + uint8_t* rows); + + // Checks each column can be a key of the kind it was before, laying the + // keys out on the first batch, then writes them. Out of line: a batch whose + // columns are shaped as the last batch's were needs none of it. + std::optional Classify(const RowBatch& batch, + const std::vector& columns); + void WriteKeys(const RowBatch& batch, const std::vector& columns); + + // What Encode() needs of one column, in 16 bytes. `shape` is its kind and + // storage type as of the last batch classified. + struct KeyColumn { + uint16_t shape; + RowLayout::Slot slot; + Writer writer; + }; + + std::vector key_columns_; + std::vector kinds_; + RowLayout layout_; + FlexVector bytes_; +}; + +} // namespace perfetto::trace_processor::core::exec + +#endif // SRC_TRACE_PROCESSOR_CORE_EXEC_KEY_ENCODER_H_ diff --git a/src/trace_processor/core/exec/key_encoder_benchmark.cc b/src/trace_processor/core/exec/key_encoder_benchmark.cc new file mode 100644 index 00000000000..40f6f1ae77b --- /dev/null +++ b/src/trace_processor/core/exec/key_encoder_benchmark.cc @@ -0,0 +1,87 @@ +/* + * Copyright (C) 2026 The Android Open Source Project + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +#include "src/trace_processor/core/exec/key_encoder.h" + +#include + +#include +#include + +#include "src/trace_processor/core/common/storage_types.h" +#include "src/trace_processor/core/exec/column_view.h" +#include "src/trace_processor/core/exec/row_batch.h" +#include "src/trace_processor/core/util/bit_vector.h" + +// Laying out keys: a full batch at a time, as an operator keyed by them does +// with its input, and a row at a time, as when a small input is run over and +// over. + +namespace perfetto::trace_processor::core::exec { +namespace { + +constexpr uint32_t kRows = kMaxBatchRows; + +struct Int64s { + Int64s() : validity(BitVector::CreateWithSize(kRows)) { + for (uint32_t i = 0; i < kRows; ++i) { + values.push_back(int64_t{i * 7919 % 1000}); + // One row in eight is null, unpredictably. + if ((i * 2654435761u) >> 29) { + validity.set(i); + } + } + } + std::vector values; + BitVector validity; +}; + +void Run(benchmark::State& state, const ColumnView& column, uint32_t rows) { + RowBatch batch; + batch.AddColumn(column); + batch.SetCardinality(rows); + std::vector columns = {0}; + KeyEncoder encoder; + for (auto _ : state) { + benchmark::DoNotOptimize(encoder.Encode(batch, columns)); + benchmark::DoNotOptimize(encoder.Key(rows - 1).data()); + } + state.SetItemsProcessed(static_cast(state.iterations()) * rows); +} + +void BM_KeyEncoderInt64(benchmark::State& state) { + Int64s c; + Run(state, ColumnView::Reference(StorageType{Int64{}}, c.values.data()), + kRows); +} +BENCHMARK(BM_KeyEncoderInt64); + +void BM_KeyEncoderInt64Nulls(benchmark::State& state) { + Int64s c; + Run(state, + ColumnView::Reference(StorageType{Int64{}}, c.values.data(), &c.validity), + kRows); +} +BENCHMARK(BM_KeyEncoderInt64Nulls); + +void BM_KeyEncoderInt64OneRow(benchmark::State& state) { + Int64s c; + Run(state, ColumnView::Reference(StorageType{Int64{}}, c.values.data()), 1); +} +BENCHMARK(BM_KeyEncoderInt64OneRow); + +} // namespace +} // namespace perfetto::trace_processor::core::exec diff --git a/src/trace_processor/core/exec/key_encoder_unittest.cc b/src/trace_processor/core/exec/key_encoder_unittest.cc new file mode 100644 index 00000000000..7697e50a65e --- /dev/null +++ b/src/trace_processor/core/exec/key_encoder_unittest.cc @@ -0,0 +1,99 @@ +/* + * Copyright (C) 2026 The Android Open Source Project + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +#include "src/trace_processor/core/exec/key_encoder.h" + +#include +#include +#include +#include + +#include "src/trace_processor/containers/string_pool.h" +#include "src/trace_processor/core/common/storage_types.h" +#include "src/trace_processor/core/exec/column_view.h" +#include "src/trace_processor/core/exec/row_batch.h" +#include "src/trace_processor/core/exec/variant.h" +#include "src/trace_processor/core/util/bit_vector.h" +#include "test/gtest_and_gmock.h" + +namespace perfetto::trace_processor::core::exec { +namespace { + +std::optional> Keys(const ColumnView& column, + uint32_t rows) { + RowBatch batch; + batch.AddColumn(column); + batch.SetCardinality(rows); + KeyEncoder encoder; + if (encoder.Encode(batch, {0})) { + return std::nullopt; + } + std::vector keys; + for (uint32_t row = 0; row < rows; ++row) { + keys.emplace_back(encoder.Key(row)); + } + return keys; +} + +TEST(KeyEncoderTest, IntegersOfEveryWidthAgree) { + int32_t i32[] = {-1, 7}; + uint32_t u32[] = {0, 7}; + int64_t i64[] = {-1, 7}; + auto a = *Keys(ColumnView::Reference(StorageType{Int32{}}, i32), 2); + auto b = *Keys(ColumnView::Reference(StorageType{Uint32{}}, u32), 2); + auto c = *Keys(ColumnView::Reference(StorageType{Int64{}}, i64), 2); + auto ids = *Keys(ColumnView::Reference(StorageType{Id{}}, nullptr), 8); + EXPECT_EQ(a[1], b[1]); + EXPECT_EQ(a[1], c[1]); + EXPECT_EQ(a[0], c[0]); + EXPECT_EQ(ids[7], c[1]); +} + +TEST(KeyEncoderTest, NullsAgreeHoweverTheyAreMarked) { + StringPool pool; + StringPool::Id strings[] = {pool.InternString("a"), StringPool::Id::Null()}; + BitVector validity = BitVector::CreateWithSize(2); + validity.set(1); + auto by_validity = *Keys( + ColumnView::Reference(StorageType{String{}}, strings, &validity), 2); + auto by_id = *Keys(ColumnView::Reference(StorageType{String{}}, strings), 2); + EXPECT_NE(by_id[0], by_id[1]); + EXPECT_EQ(by_validity[0], by_id[1]); + // The null id is 0, so a string needs no null byte. + EXPECT_EQ(by_id[0].size(), sizeof(uint32_t)); +} + +TEST(KeyEncoderTest, OnlyColumnsOfOneTypeAreKeys) { + double doubles[] = {1.5}; + EXPECT_TRUE(Keys(ColumnView::Reference(StorageType{Double{}}, doubles), 1)); + Variant variants[] = {Variant{}}; + EXPECT_FALSE(Keys(ColumnView::Variants(variants), 1)); +} + +TEST(KeyEncoderTest, AColumnKeepsItsType) { + int64_t ints[] = {1}; + double doubles[] = {1.0}; + RowBatch batch; + batch.AddColumn(ColumnView::Reference(StorageType{Int64{}}, ints)); + batch.AddColumn(ColumnView::Reference(StorageType{Double{}}, doubles)); + batch.SetCardinality(1); + KeyEncoder encoder; + EXPECT_EQ(encoder.Encode(batch, {0}), std::nullopt); + EXPECT_EQ(encoder.Encode(batch, {1}), 0u); +} + +} // namespace +} // namespace perfetto::trace_processor::core::exec diff --git a/src/trace_processor/perfetto_sql/engine/perfetto_sql_connection_unittest.cc b/src/trace_processor/perfetto_sql/engine/perfetto_sql_connection_unittest.cc index ee06d4d49ad..3d73bd1bcfa 100644 --- a/src/trace_processor/perfetto_sql/engine/perfetto_sql_connection_unittest.cc +++ b/src/trace_processor/perfetto_sql/engine/perfetto_sql_connection_unittest.cc @@ -1175,6 +1175,33 @@ TEST_F(PerfettoSqlConnectionPipelineTest, NewDataframeArgsAreSafe) { EXPECT_EQ(rows->size(), 3u); } +// Nulls agree with each other, as under GROUP BY. +TEST_F(PerfettoSqlConnectionPipelineTest, IntervalIntersectionPerAnyType) { + ASSERT_TRUE(Rows(R"( + CREATE TABLE a(ts INTEGER, dur INTEGER, name TEXT, n INTEGER); + CREATE TABLE b(ts INTEGER, dur INTEGER, name TEXT, n INTEGER); + INSERT INTO a VALUES (0, 10, 'x', 1), (0, 10, 'y', 2), (0, 10, NULL, 3); + INSERT INTO b VALUES (5, 10, 'x', 1), (5, 10, NULL, 2), (5, 10, 'z', 3); + )") + .ok()); + auto rows = Rows(R"( + INTERVAL INTERSECTION OF (a AS a, b AS b) PER name + |> SELECT ts, dur, name + )"); + ASSERT_TRUE(rows.ok()) << rows.status().message(); + EXPECT_THAT(*rows, testing::ElementsAre("5,5,NULL", "5,5,x")); + + // Keys must be the same type in every operand. + EXPECT_THAT(Rows(R"( + INTERVAL INTERSECTION OF ( + a AS a, (SELECT ts, dur, n AS name FROM b) AS b + ) PER name + )") + .status() + .message(), + testing::HasSubstr("the same in every operand")); +} + TEST_F(PerfettoSqlConnectionPipelineTest, ForksRunPipelinesIndependently) { auto fork = connection_->Fork(); auto first = connection_->ExecuteUntilLastStatement( diff --git a/src/trace_processor/perfetto_sql/pipeline/physical_plan.cc b/src/trace_processor/perfetto_sql/pipeline/physical_plan.cc index 8b9e835bf0a..7ceeefc8d53 100644 --- a/src/trace_processor/perfetto_sql/pipeline/physical_plan.cc +++ b/src/trace_processor/perfetto_sql/pipeline/physical_plan.cc @@ -169,8 +169,8 @@ void Lowering::LowerIntervalIntersect(const op::IntervalIntersect& isect, } return at; }; - // An operand is read through a pipeline of its own, which widens the - // columns the intersection reads to Int64 where they are not already. + // An operand is read through a pipeline of its own, which widens its + // bounds to Int64 where they are not already. std::vector> widen; ex::IntervalIntersectOperand lowered; lowered.ts_column = position(operand.ts); @@ -183,9 +183,6 @@ void Lowering::LowerIntervalIntersect(const op::IntervalIntersect& isect, } RequireInt64(operand.ts, lowered.ts_column, widen); RequireInt64(operand.dur, lowered.dur_column, widen); - for (uint32_t k = 0; k < operand.keys.size(); k++) { - RequireInt64(operand.keys[k], lowered.key_columns[k], widen); - } out_->operand_inputs_.push_back(MakeSource(scan)); out_->operand_pipelines_.push_back(std::make_unique( *out_->operand_inputs_.back(), std::move(widen), From 7498915fb74d21cfe1f8438f2b7770d05952ee72 Mon Sep 17 00:00:00 2001 From: Lalit Maganti Date: Fri, 2 Oct 2026 16:12:16 +0100 Subject: [PATCH 3/3] tp: share checked flat-column reads in key encoding --- src/trace_processor/core/exec/column_view.h | 5 +- src/trace_processor/core/exec/key_encoder.cc | 46 +++++++++---------- .../core/exec/key_encoder_unittest.cc | 42 +++++++++++++++++ 3 files changed, 69 insertions(+), 24 deletions(-) diff --git a/src/trace_processor/core/exec/column_view.h b/src/trace_processor/core/exec/column_view.h index f69130437e3..d18b841fcd0 100644 --- a/src/trace_processor/core/exec/column_view.h +++ b/src/trace_processor/core/exec/column_view.h @@ -138,7 +138,10 @@ class FlatColumnReader { explicit FlatColumnReader(const ColumnView& column) : data_(static_cast(column.data())), selection_(column.selection()), - validity_(column.validity()) {} + validity_(column.validity()) { + PERFETTO_DCHECK(column.kind() == ColumnView::Kind::kFlat); + PERFETTO_DCHECK(column.type().Is::type>()); + } // False if the row holds no value. PERFETTO_ALWAYS_INLINE bool Read(uint32_t row, T* out) const { diff --git a/src/trace_processor/core/exec/key_encoder.cc b/src/trace_processor/core/exec/key_encoder.cc index f1b2faec079..128aa9cc57e 100644 --- a/src/trace_processor/core/exec/key_encoder.cc +++ b/src/trace_processor/core/exec/key_encoder.cc @@ -34,52 +34,52 @@ namespace perfetto::trace_processor::core::exec { namespace { -template -PERFETTO_ALWAYS_INLINE void Write(const ColumnView& column, - RowLayout::Slot slot, - uint32_t count, - uint8_t* rows, - Value value) { +void WriteSequence(const ColumnView& column, + RowLayout::Slot slot, + uint32_t count, + uint8_t* rows) { RowSelection selection = column.selection(); const BitVector* validity = column.validity(); - RowLayout::Write( + RowLayout::Write( slot, count, - [&](uint32_t row, T* out) { + [&](uint32_t row, int64_t* out) { uint32_t index = selection.GetIndex(row); if (validity && !validity->is_set(index)) { return false; } - *out = value(index); + *out = int64_t{index}; return true; }, rows); } -void WriteSequence(const ColumnView& column, - RowLayout::Slot slot, - uint32_t count, - uint8_t* rows) { - Write(column, slot, count, rows, - [](uint32_t index) { return int64_t{index}; }); -} - template void WriteInteger(const ColumnView& column, RowLayout::Slot slot, uint32_t count, uint8_t* rows) { - const auto* data = static_cast(column.data()); - Write(column, slot, count, rows, - [data](uint32_t index) { return int64_t{data[index]}; }); + FlatColumnReader reader(column); + RowLayout::Write( + slot, count, + [&](uint32_t row, int64_t* out) { + Stored value; + if (!reader.Read(row, &value)) { + return false; + } + *out = int64_t{value}; + return true; + }, + rows); } void WriteDouble(const ColumnView& column, RowLayout::Slot slot, uint32_t count, uint8_t* rows) { - const auto* data = static_cast(column.data()); - Write(column, slot, count, rows, - [data](uint32_t index) { return data[index]; }); + FlatColumnReader reader(column); + RowLayout::Write( + slot, count, + [&](uint32_t row, double* out) { return reader.Read(row, out); }, rows); } // Not nullable: the null id is 0. diff --git a/src/trace_processor/core/exec/key_encoder_unittest.cc b/src/trace_processor/core/exec/key_encoder_unittest.cc index 7697e50a65e..f6b24c5cf9b 100644 --- a/src/trace_processor/core/exec/key_encoder_unittest.cc +++ b/src/trace_processor/core/exec/key_encoder_unittest.cc @@ -25,6 +25,7 @@ #include "src/trace_processor/core/common/storage_types.h" #include "src/trace_processor/core/exec/column_view.h" #include "src/trace_processor/core/exec/row_batch.h" +#include "src/trace_processor/core/exec/row_selection.h" #include "src/trace_processor/core/exec/variant.h" #include "src/trace_processor/core/util/bit_vector.h" #include "test/gtest_and_gmock.h" @@ -76,6 +77,47 @@ TEST(KeyEncoderTest, NullsAgreeHoweverTheyAreMarked) { EXPECT_EQ(by_id[0].size(), sizeof(uint32_t)); } +TEST(KeyEncoderTest, SelectedRowsKeepValuesAndNulls) { + BitVector validity = BitVector::CreateWithSize(3); + validity.set(0); + validity.set(2); + auto check = [](ColumnView column) { + auto original = Keys(column, 3); + ASSERT_TRUE(original.has_value()); + EXPECT_NE((*original)[0], (*original)[1]); + EXPECT_NE((*original)[0], (*original)[2]); + + column.SetRange(1); + auto ranged = Keys(column, 2); + ASSERT_TRUE(ranged.has_value()); + EXPECT_EQ(*ranged, + (std::vector{(*original)[1], (*original)[2]})); + + // Compose reordered, repeated indices with the range. Validity must be + // checked at the physical row, not the logical row in this selection. + uint32_t indices[] = {1, 0, 1}; + SelectionPool pool; + column.Slice(RowSelection::Indices(indices), 3, pool); + auto selected = Keys(column, 3); + ASSERT_TRUE(selected.has_value()); + EXPECT_EQ(*selected, (std::vector{ + (*original)[2], (*original)[1], (*original)[2]})); + }; + int32_t i32[] = {-1, 99, 7}; + uint32_t u32[] = {3, 99, 7}; + int64_t i64[] = {-1, 99, 7}; + double doubles[] = {1.5, 99.0, 7.5}; + StringPool pool; + StringPool::Id strings[] = {pool.InternString("a"), pool.InternString("b"), + pool.InternString("c")}; + check(ColumnView::Reference(StorageType{Int32{}}, i32, &validity)); + check(ColumnView::Reference(StorageType{Uint32{}}, u32, &validity)); + check(ColumnView::Reference(StorageType{Int64{}}, i64, &validity)); + check(ColumnView::Reference(StorageType{Double{}}, doubles, &validity)); + check(ColumnView::Reference(StorageType{String{}}, strings, &validity)); + check(ColumnView::Reference(StorageType{Id{}}, nullptr, &validity)); +} + TEST(KeyEncoderTest, OnlyColumnsOfOneTypeAreKeys) { double doubles[] = {1.5}; EXPECT_TRUE(Keys(ColumnView::Reference(StorageType{Double{}}, doubles), 1));