Skip to content

Commit f40a71e

Browse files
author
wangyong.alen
committed
feat(scan): prune append buckets from equality predicates
1 parent ead9207 commit f40a71e

5 files changed

Lines changed: 228 additions & 0 deletions

File tree

‎docs/source/api/scan.rst‎

Lines changed: 11 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -21,6 +21,17 @@ Scan
2121

2222
.. _cpp-api-scan:
2323

24+
Bucket pruning
25+
==============
26+
27+
For fixed-bucket append tables, an equality predicate on every bucket key lets
28+
the scan derive the target bucket using the table's bucket function. Other buckets
29+
are excluded from the scan plan without requiring an explicit bucket ID from the
30+
caller. An explicit bucket filter takes precedence. Queries that do not constrain
31+
all bucket keys with equality, and bucket-unaware tables, keep the existing scan
32+
behavior. Inferred pruning applies only to files matching the scan schema and
33+
bucket count; older layouts retain the existing filtering behavior.
34+
2435
Interface
2536
=========
2637

‎src/paimon/core/operation/append_only_file_store_scan.cpp‎

Lines changed: 24 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -26,11 +26,13 @@
2626
#include <set>
2727
#include <string>
2828
#include <utility>
29+
#include <vector>
2930

3031
#include "arrow/type.h"
3132
#include "fmt/format.h"
3233
#include "paimon/common/predicate/predicate_filter.h"
3334
#include "paimon/common/types/data_field.h"
35+
#include "paimon/core/bucket/bucket_select_converter.h"
3436
#include "paimon/core/core_options.h"
3537
#include "paimon/core/io/data_file_meta.h"
3638
#include "paimon/core/io/file_index_evaluator.h"
@@ -42,6 +44,7 @@
4244
#include "paimon/core/utils/field_mapping.h"
4345
#include "paimon/file_index/file_index_result.h"
4446
#include "paimon/predicate/predicate_utils.h"
47+
#include "paimon/scan_context.h"
4548
#include "paimon/status.h"
4649

4750
namespace paimon {
@@ -66,6 +69,21 @@ Result<std::unique_ptr<AppendOnlyFileStoreScan>> AppendOnlyFileStoreScan::Create
6669
table_schema, arrow_schema, core_options, executor, pool));
6770
PAIMON_RETURN_NOT_OK(
6871
scan->SplitAndSetFilter(table_schema->PartitionKeys(), arrow_schema, scan_filters));
72+
const auto& bucket_keys = table_schema->BucketKeys();
73+
int32_t num_buckets = core_options.GetBucket();
74+
if (scan->predicates_ && !scan_filters->GetBucketFilter().has_value() && num_buckets > 0 &&
75+
!bucket_keys.empty()) {
76+
std::vector<std::shared_ptr<arrow::DataType>> bucket_key_types;
77+
bucket_key_types.reserve(bucket_keys.size());
78+
for (const auto& key : bucket_keys) {
79+
PAIMON_ASSIGN_OR_RAISE(DataField field, table_schema->GetField(key));
80+
bucket_key_types.push_back(field.Type());
81+
}
82+
PAIMON_ASSIGN_OR_RAISE(scan->predicate_bucket_,
83+
BucketSelectConverter::Convert(
84+
scan->predicates_, bucket_keys, bucket_key_types,
85+
core_options.GetBucketFunctionType(), num_buckets, pool.get()));
86+
}
6987
return scan;
7088
}
7189

@@ -86,6 +104,12 @@ Result<bool> AppendOnlyFileStoreScan::FilterByStats(const ManifestEntry& entry)
86104
if (!predicates_) {
87105
return true;
88106
}
107+
// A historical file may use a different schema or bucket count after a rescale.
108+
// Keep the inferred bucket separate from the caller's explicit bucket filter.
109+
if (predicate_bucket_ && entry.TotalBuckets() == core_options_.GetBucket() &&
110+
entry.File()->schema_id == table_schema_->Id() && entry.Bucket() != *predicate_bucket_) {
111+
return false;
112+
}
89113
const auto& meta = entry.File();
90114
std::shared_ptr<TableSchema> data_schema = table_schema_;
91115
std::shared_ptr<Predicate> trimmed_predicates = predicates_;

‎src/paimon/core/operation/append_only_file_store_scan.h‎

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -18,7 +18,9 @@
1818

1919
#pragma once
2020

21+
#include <cstdint>
2122
#include <memory>
23+
#include <optional>
2224
#include <vector>
2325

2426
#include "paimon/core/operation/file_store_scan.h"
@@ -79,5 +81,6 @@ class AppendOnlyFileStoreScan : public FileStoreScan {
7981

8082
private:
8183
std::shared_ptr<SimpleStatsEvolutions> evolutions_;
84+
std::optional<int32_t> predicate_bucket_;
8285
};
8386
} // namespace paimon

‎src/paimon/core/operation/append_only_file_store_scan_test.cpp‎

Lines changed: 99 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -20,15 +20,18 @@
2020

2121
#include <algorithm>
2222
#include <cstdint>
23+
#include <map>
2324
#include <optional>
2425
#include <string>
2526
#include <utility>
2627
#include <vector>
2728

29+
#include "arrow/type.h"
2830
#include "gtest/gtest.h"
2931
#include "paimon/common/data/binary_row.h"
3032
#include "paimon/common/data/binary_row_writer.h"
3133
#include "paimon/common/io/cache/lru_cache.h"
34+
#include "paimon/core/bucket/default_bucket_function.h"
3235
#include "paimon/core/manifest/manifest_entry.h"
3336
#include "paimon/core/manifest/partition_entry.h"
3437
#include "paimon/core/schema/schema_manager.h"
@@ -37,6 +40,7 @@
3740
#include "paimon/core/table/source/abstract_table_scan.h"
3841
#include "paimon/core/table/source/snapshot/snapshot_reader.h"
3942
#include "paimon/defs.h"
43+
#include "paimon/executor.h"
4044
#include "paimon/fs/local/local_file_system.h"
4145
#include "paimon/memory/memory_pool.h"
4246
#include "paimon/metrics.h"
@@ -46,10 +50,105 @@
4650
#include "paimon/status.h"
4751
#include "paimon/table/source/scan_metrics.h"
4852
#include "paimon/table/source/table_scan.h"
53+
#include "paimon/testing/utils/binary_row_generator.h"
4954
#include "paimon/testing/utils/testharness.h"
5055
#include "paimon/testing/utils/timezone_guard.h"
5156
namespace paimon::test {
5257

58+
class AppendBucketPruningTest : public testing::Test {
59+
public:
60+
Result<std::unique_ptr<AppendOnlyFileStoreScan>> CreateScan(
61+
const std::shared_ptr<Predicate>& predicate,
62+
const std::optional<int32_t>& bucket = std::nullopt) const {
63+
auto arrow_schema = DataField::ConvertDataFieldsToArrowSchema(
64+
{DataField(0, arrow::field("rowkey", arrow::utf8())),
65+
DataField(1, arrow::field("value", arrow::int32()))});
66+
PAIMON_ASSIGN_OR_RAISE(std::shared_ptr<TableSchema> schema,
67+
TableSchema::Create(schema_id_, arrow_schema, {}, {}, options_));
68+
PAIMON_ASSIGN_OR_RAISE(CoreOptions options, CoreOptions::FromMap(options_));
69+
auto filters = std::make_shared<ScanFilter>(
70+
predicate, std::vector<std::map<std::string, std::string>>(), bucket);
71+
return AppendOnlyFileStoreScan::Create(nullptr, schema_manager_, nullptr, nullptr, schema,
72+
arrow_schema, filters, options,
73+
GetGlobalDefaultExecutor(), pool_);
74+
}
75+
76+
void CheckBuckets(const std::shared_ptr<Predicate>& predicate,
77+
const std::optional<int32_t>& expected_bucket,
78+
const std::optional<int32_t>& explicit_bucket = std::nullopt,
79+
int32_t total_buckets = kNumBuckets) const {
80+
ASSERT_OK_AND_ASSIGN(auto scan, CreateScan(predicate, explicit_bucket));
81+
// Every file's stats match the lookup, so only bucket pruning can discard a file.
82+
SimpleStats stats = BinaryRowGenerator::GenerateStats(
83+
{std::string("a"), 0}, {std::string("z"), 100}, {0, 0}, pool_.get());
84+
ASSERT_OK_AND_ASSIGN(
85+
auto file,
86+
DataFileMeta::ForAppend("data.parquet", 100, 10, stats, 0, 9, 0, std::nullopt,
87+
std::nullopt, std::nullopt, std::nullopt, std::nullopt));
88+
for (int32_t bucket = 0; bucket < total_buckets; ++bucket) {
89+
ManifestEntry entry(FileKind::Add(), BinaryRow::EmptyRow(), bucket, total_buckets,
90+
file);
91+
ASSERT_OK_AND_ASSIGN(auto stats_scan, CreateScan(predicate, bucket));
92+
ASSERT_OK_AND_ASSIGN(bool stats_match, stats_scan->FilterByStats(entry));
93+
ASSERT_TRUE(stats_match);
94+
ASSERT_OK_AND_ASSIGN(bool keep, scan->FilterManifestEntry(entry));
95+
ASSERT_EQ(keep, !expected_bucket.has_value() || bucket == expected_bucket.value());
96+
}
97+
}
98+
99+
std::shared_ptr<Predicate> KeyEquals() const {
100+
return PredicateBuilder::Equal(0, "rowkey", FieldType::STRING,
101+
Literal(FieldType::STRING, "key", 3));
102+
}
103+
104+
static constexpr int32_t kNumBuckets = 4;
105+
int64_t schema_id_ = 0;
106+
std::shared_ptr<SchemaManager> schema_manager_;
107+
std::shared_ptr<MemoryPool> pool_ = GetDefaultPool();
108+
std::map<std::string, std::string> options_ = {{Options::BUCKET, "4"},
109+
{Options::BUCKET_KEY, "rowkey"}};
110+
};
111+
112+
TEST_F(AppendBucketPruningTest, PrunesStringKeyLookup) {
113+
BinaryRow key = BinaryRowGenerator::GenerateRow({std::string("key")}, pool_.get());
114+
int32_t expected_bucket = DefaultBucketFunction().Bucket(key, kNumBuckets);
115+
CheckBuckets(KeyEquals(), expected_bucket);
116+
}
117+
118+
TEST_F(AppendBucketPruningTest, PreservesExplicitBucketFilter) {
119+
BinaryRow key = BinaryRowGenerator::GenerateRow({std::string("key")}, pool_.get());
120+
int32_t other_bucket = (DefaultBucketFunction().Bucket(key, kNumBuckets) + 1) % kNumBuckets;
121+
CheckBuckets(KeyEquals(), other_bucket, other_bucket);
122+
}
123+
124+
TEST_F(AppendBucketPruningTest, DoesNotPruneWithoutCompleteEqualKeys) {
125+
CheckBuckets(nullptr, std::nullopt);
126+
CheckBuckets(PredicateBuilder::GreaterThan(0, "rowkey", FieldType::STRING,
127+
Literal(FieldType::STRING, "key", 3)),
128+
std::nullopt);
129+
options_[Options::BUCKET_KEY] = "rowkey,value";
130+
CheckBuckets(KeyEquals(), std::nullopt);
131+
}
132+
133+
TEST_F(AppendBucketPruningTest, DoesNotPruneBucketUnawareTable) {
134+
options_[Options::BUCKET] = "-1";
135+
CheckBuckets(KeyEquals(), std::nullopt);
136+
}
137+
138+
TEST_F(AppendBucketPruningTest, DoesNotPruneDifferentBucketCount) {
139+
CheckBuckets(KeyEquals(), std::nullopt, std::nullopt, 8);
140+
}
141+
142+
TEST_F(AppendBucketPruningTest, DoesNotPruneDifferentSchema) {
143+
auto test_dir = UniqueTestDirectory::Create("local");
144+
schema_manager_ =
145+
std::make_shared<SchemaManager>(std::make_shared<LocalFileSystem>(), test_dir->Str());
146+
ASSERT_OK_AND_ASSIGN(auto scan, CreateScan(nullptr));
147+
ASSERT_OK(schema_manager_->CreateTable(scan->schema_, {}, {}, options_));
148+
schema_id_ = 1;
149+
CheckBuckets(KeyEquals(), std::nullopt);
150+
}
151+
53152
TEST(AppendOnlyFileStoreScanTest, TestReconstructPredicateWithNonCastedFields) {
54153
std::string table_root =
55154
paimon::test::GetDataDir() +

‎test/inte/scan_and_read_inte_test.cpp‎

Lines changed: 91 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -29,8 +29,12 @@
2929
#include <vector>
3030

3131
#include "arrow/api.h"
32+
#include "arrow/c/bridge.h"
3233
#include "arrow/ipc/json_simple.h"
34+
#include "fmt/format.h"
35+
#include "fmt/ranges.h"
3336
#include "gtest/gtest.h"
37+
#include "paimon/bucket/bucket_id_calculator.h"
3438
#include "paimon/common/factories/io_hook.h"
3539
#include "paimon/common/io/cache/lru_cache.h"
3640
#include "paimon/common/table/special_fields.h"
@@ -421,6 +425,93 @@ TEST_P(ScanAndReadInteTest, TestWithAppendSnapshot3) {
421425
ASSERT_TRUE(expected->Equals(read_result)) << read_result->ToString();
422426
}
423427

428+
TEST_P(ScanAndReadInteTest, TestWithAppendBucketKeyPointLookup) {
429+
for (const auto& key_type : {arrow::utf8(), arrow::binary()}) {
430+
SCOPED_TRACE(key_type->ToString());
431+
auto test_dir = UniqueTestDirectory::Create("local");
432+
arrow::FieldVector fields = {arrow::field("rowkey", key_type),
433+
arrow::field("value", arrow::int32())};
434+
std::map<std::string, std::string> options = {{Options::FILE_FORMAT, FileFormat()},
435+
{Options::FILE_SYSTEM, "local"},
436+
{Options::BUCKET, "4"},
437+
{Options::BUCKET_KEY, "rowkey"}};
438+
ASSERT_OK_AND_ASSIGN(
439+
auto helper, TestHelper::Create(test_dir->Str(), arrow::schema(fields), {}, {}, options,
440+
/*is_streaming_mode=*/false));
441+
std::string table_path = test_dir->Str() + "/foo.db/bar";
442+
// Route writes through the public bucket calculator, as a KV client does.
443+
constexpr int32_t kNumKeys = 64;
444+
constexpr int32_t kNumBuckets = 4;
445+
std::vector<std::string> keys;
446+
for (int32_t i = 0; i < kNumKeys; ++i) {
447+
keys.push_back(fmt::format(R"(["key{:03}"])", i));
448+
}
449+
auto key_array = arrow::ipc::internal::json::ArrayFromJSON(
450+
arrow::struct_({fields[0]}), fmt::format("[{}]", fmt::join(keys, ",")))
451+
.ValueOrDie();
452+
::ArrowArray c_keys;
453+
::ArrowSchema c_schema;
454+
ASSERT_TRUE(arrow::ExportArray(*key_array, &c_keys, &c_schema).ok());
455+
ASSERT_OK_AND_ASSIGN(auto calculator,
456+
BucketIdCalculator::Create(false, kNumBuckets, GetDefaultPool()));
457+
std::vector<int32_t> bucket_ids(kNumKeys);
458+
ASSERT_OK(calculator->CalculateBucketIds(&c_keys, &c_schema, bucket_ids.data()));
459+
std::vector<std::vector<std::string>> rows(kNumBuckets);
460+
for (int32_t i = 0; i < kNumKeys; ++i) {
461+
rows[bucket_ids[i]].push_back(fmt::format(R"(["key{:03}", {}])", i, i));
462+
}
463+
std::vector<std::unique_ptr<RecordBatch>> batches;
464+
for (int32_t bucket = 0; bucket < kNumBuckets; ++bucket) {
465+
ASSERT_FALSE(rows[bucket].empty());
466+
ASSERT_OK_AND_ASSIGN(
467+
auto batch, TestHelper::MakeRecordBatch(
468+
arrow::struct_(fields),
469+
fmt::format("[{}]", fmt::join(rows[bucket], ",")), {}, bucket, {}));
470+
batches.push_back(std::move(batch));
471+
}
472+
ASSERT_OK(helper->WriteAndCommit(std::move(batches), 0, std::nullopt));
473+
474+
FieldType field_type =
475+
key_type->id() == arrow::Type::STRING ? FieldType::STRING : FieldType::BINARY;
476+
auto predicate =
477+
PredicateBuilder::Equal(0, "rowkey", field_type, Literal(field_type, "key032", 6));
478+
// All buckets have overlapping stats for this key. An explicit filter bypasses
479+
// inference, proving that ordinary statistics cannot account for the pruning.
480+
for (int32_t bucket = 0; bucket < kNumBuckets; ++bucket) {
481+
ScanContextBuilder builder(table_path);
482+
builder.SetPredicate(predicate).SetBucketFilter(bucket);
483+
ASSERT_OK_AND_ASSIGN(auto context, FinishScanContext(builder));
484+
ASSERT_OK_AND_ASSIGN(auto scan, TableScan::Create(std::move(context)));
485+
ASSERT_OK_AND_ASSIGN(auto plan, scan->CreatePlan());
486+
ASSERT_FALSE(plan->Splits().empty());
487+
}
488+
ScanContextBuilder scan_builder(table_path);
489+
scan_builder.SetPredicate(predicate);
490+
ASSERT_OK_AND_ASSIGN(auto scan_context, FinishScanContext(scan_builder));
491+
ASSERT_OK_AND_ASSIGN(auto scan, TableScan::Create(std::move(scan_context)));
492+
ASSERT_OK_AND_ASSIGN(auto plan, scan->CreatePlan());
493+
ASSERT_FALSE(plan->Splits().empty());
494+
for (const auto& split : plan->Splits()) {
495+
auto data_split = std::dynamic_pointer_cast<DataSplitImpl>(split);
496+
ASSERT_TRUE(data_split);
497+
ASSERT_EQ(data_split->Bucket(), bucket_ids[32]);
498+
}
499+
ReadContextBuilder read_builder(table_path);
500+
AddReadOptionsForPrefetch(&read_builder);
501+
read_builder.SetPredicate(predicate).EnablePredicateFilter(true);
502+
ASSERT_OK_AND_ASSIGN(auto read_context, read_builder.Finish());
503+
ASSERT_OK_AND_ASSIGN(auto read, TableRead::Create(std::move(read_context)));
504+
ASSERT_OK_AND_ASSIGN(auto reader, read->CreateReader(plan->Splits()));
505+
ASSERT_OK_AND_ASSIGN(auto result, ReadResultCollector::CollectResult(std::move(reader)));
506+
fields.insert(fields.begin(), arrow::field("_VALUE_KIND", arrow::int8()));
507+
auto expected_array = arrow::ipc::internal::json::ArrayFromJSON(arrow::struct_(fields),
508+
R"([[0, "key032", 32]])")
509+
.ValueOrDie();
510+
auto expected = std::make_shared<arrow::ChunkedArray>(expected_array);
511+
ASSERT_TRUE(expected->Equals(result)) << result->ToString();
512+
}
513+
}
514+
424515
TEST_P(ScanAndReadInteTest, TestWithAppendSnapshot5) {
425516
auto file_format = FileFormat();
426517
std::string table_path = GetDataDir() + "/" + file_format + "/append_09.db/append_09";

0 commit comments

Comments
 (0)