Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
27 changes: 27 additions & 0 deletions include/paimon/scan_context.h
Original file line number Diff line number Diff line change
Expand Up @@ -37,6 +37,7 @@ class ScanFilter;
class Executor;
class MemoryPool;
class Predicate;
class SnapshotReadView;

/// `ScanContext` is some configuration for table scan operations.
///
Expand All @@ -56,6 +57,18 @@ class PAIMON_EXPORT ScanContext {
const std::map<std::string, std::string>& options,
const std::shared_ptr<Cache>& cache);

ScanContext(const std::string& path, bool is_streaming_mode, std::optional<int32_t> limit,
const std::shared_ptr<ScanFilter>& scan_filter,
const std::shared_ptr<GlobalIndexResult>& global_index_result,
const std::shared_ptr<RealtimeContext>& realtime_context,
const std::shared_ptr<MemoryPool>& memory_pool,
const std::shared_ptr<Executor>& executor,
const std::shared_ptr<FileSystem>& specific_file_system,
const std::optional<std::string>& table_schema,
const std::map<std::string, std::string>& options,
const std::shared_ptr<Cache>& cache,
const std::shared_ptr<const SnapshotReadView>& snapshot_read_view);

~ScanContext();

const std::string& GetPath() const {
Expand Down Expand Up @@ -105,6 +118,10 @@ class PAIMON_EXPORT ScanContext {
return cache_;
}

std::shared_ptr<const SnapshotReadView> GetSnapshotReadView() const {
return snapshot_read_view_;
}

private:
std::string path_;
bool is_streaming_mode_;
Expand All @@ -118,6 +135,7 @@ class PAIMON_EXPORT ScanContext {
std::optional<std::string> table_schema_;
std::map<std::string, std::string> options_;
std::shared_ptr<Cache> cache_;
std::shared_ptr<const SnapshotReadView> snapshot_read_view_;
};

/// Filter configuration for table scan operations
Expand Down Expand Up @@ -220,6 +238,15 @@ class PAIMON_EXPORT ScanContextBuilder {
/// @return Reference to this builder for method chaining.
ScanContextBuilder& WithCache(const std::shared_ptr<Cache>& cache);

/// Reuse the immutable parsed snapshot metadata from an earlier batch scan plan.
///
/// The view must belong to the same normalized table path and branch. Snapshot read views are
/// supported for non-streaming latest-snapshot scans.
/// @param snapshot_read_view A read view returned by `Plan::GetSnapshotReadView()`.
/// @return Reference to this builder for method chaining.
ScanContextBuilder& WithSnapshotReadView(
const std::shared_ptr<const SnapshotReadView>& snapshot_read_view);

/// Build and return a `ScanContext` instance with input validation.
/// @return Result containing the constructed `ScanContext` or an error status.
Result<std::unique_ptr<ScanContext>> Finish();
Expand Down
9 changes: 9 additions & 0 deletions include/paimon/table/source/plan.h
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,7 @@
#include <optional>
#include <vector>

#include "paimon/table/source/snapshot_read_view.h"
#include "paimon/table/source/split.h"

namespace paimon {
Expand All @@ -34,5 +35,13 @@ class PAIMON_EXPORT Plan {
virtual const std::vector<std::shared_ptr<Split>>& Splits() const = 0;
/// Snapshot id of this plan, return `std::nullopt` if the table is empty.
virtual std::optional<int64_t> SnapshotId() const = 0;
/// Immutable snapshot metadata used to build this plan.
///
/// The returned view can be injected into a later batch scan to reuse the already resolved and
/// parsed snapshot. Only non-streaming, non-real-time latest-snapshot plans publish a reusable
/// view; other plans return `nullptr`.
virtual std::shared_ptr<const SnapshotReadView> GetSnapshotReadView() const {
return nullptr;
}
Comment thread
wangyong9999 marked this conversation as resolved.
};
} // namespace paimon
61 changes: 61 additions & 0 deletions include/paimon/table/source/snapshot_read_view.h
Original file line number Diff line number Diff line change
@@ -0,0 +1,61 @@
/*
* Licensed to the Apache Software Foundation (ASF) under one
* or more contributor license agreements. See the NOTICE file
* distributed with this work for additional information
* regarding copyright ownership. The ASF licenses this file
* to you 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.
*/

#pragma once

#include <cstdint>
#include <optional>
#include <string>

#include "paimon/visibility.h"

namespace paimon {

class SnapshotReadViewImpl;

/// An immutable, process-local view of the snapshot and table-schema metadata used to create a
/// scan plan.
///
/// A view is bound to one base table path and branch. Pass a view returned by `Plan` to
/// `ScanContextBuilder::WithSnapshotReadView()` to plan another batch scan against the same
/// parsed planning metadata without resolving the latest snapshot or schema again. Views are
/// published and accepted only for non-streaming, non-real-time latest-snapshot scans. A missing
/// snapshot id represents an empty table at the time the view was created. The base path identity
/// is also used by supported system-table scans such as `$ro`, because they consume the same
/// snapshot and schema metadata.
class PAIMON_EXPORT SnapshotReadView {
public:
virtual ~SnapshotReadView() = default;

/// Normalized base table path this view belongs to, without a system-table suffix.
virtual const std::string& TablePath() const = 0;

/// Normalized branch this view belongs to.
virtual const std::string& Branch() const = 0;

/// Snapshot id captured by this view, or `std::nullopt` if the table was empty.
virtual std::optional<int64_t> SnapshotId() const = 0;

private:
SnapshotReadView() = default;

friend class SnapshotReadViewImpl;
};

} // namespace paimon
29 changes: 27 additions & 2 deletions src/paimon/core/operation/scan_context.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -39,6 +39,22 @@ ScanContext::ScanContext(const std::string& path, bool is_streaming_mode,
const std::optional<std::string>& table_schema,
const std::map<std::string, std::string>& options,
const std::shared_ptr<Cache>& cache)
: ScanContext(path, is_streaming_mode, limit, scan_filter, global_index_result,
realtime_context, memory_pool, executor, specific_file_system, table_schema,
options, cache, nullptr) {}

ScanContext::ScanContext(const std::string& path, bool is_streaming_mode,
std::optional<int32_t> limit,
const std::shared_ptr<ScanFilter>& scan_filter,
const std::shared_ptr<GlobalIndexResult>& global_index_result,
const std::shared_ptr<RealtimeContext>& realtime_context,
const std::shared_ptr<MemoryPool>& memory_pool,
const std::shared_ptr<Executor>& executor,
const std::shared_ptr<FileSystem>& specific_file_system,
const std::optional<std::string>& table_schema,
const std::map<std::string, std::string>& options,
const std::shared_ptr<Cache>& cache,
const std::shared_ptr<const SnapshotReadView>& snapshot_read_view)
: path_(path),
is_streaming_mode_(is_streaming_mode),
limit_(limit),
Expand All @@ -50,7 +66,8 @@ ScanContext::ScanContext(const std::string& path, bool is_streaming_mode,
specific_file_system_(specific_file_system),
table_schema_(table_schema),
options_(options),
cache_(cache) {}
cache_(cache),
snapshot_read_view_(snapshot_read_view) {}

ScanContext::~ScanContext() = default;

Expand All @@ -72,6 +89,7 @@ class ScanContextBuilder::Impl {
table_schema_ = std::nullopt;
options_.clear();
cache_.reset();
snapshot_read_view_.reset();
}

private:
Expand All @@ -89,6 +107,7 @@ class ScanContextBuilder::Impl {
std::optional<std::string> table_schema_;
std::map<std::string, std::string> options_;
std::shared_ptr<Cache> cache_;
std::shared_ptr<const SnapshotReadView> snapshot_read_view_;
};

ScanContextBuilder::ScanContextBuilder(const std::string& path)
Expand Down Expand Up @@ -173,6 +192,12 @@ ScanContextBuilder& ScanContextBuilder::WithCache(const std::shared_ptr<Cache>&
return *this;
}

ScanContextBuilder& ScanContextBuilder::WithSnapshotReadView(
const std::shared_ptr<const SnapshotReadView>& snapshot_read_view) {
impl_->snapshot_read_view_ = snapshot_read_view;
return *this;
}

Result<std::unique_ptr<ScanContext>> ScanContextBuilder::Finish() {
PAIMON_ASSIGN_OR_RAISE(impl_->path_, PathUtil::NormalizePath(impl_->path_));
if (impl_->path_.empty()) {
Expand All @@ -184,7 +209,7 @@ Result<std::unique_ptr<ScanContext>> ScanContextBuilder::Finish() {
impl_->bucket_filter_),
impl_->global_index_result_, impl_->realtime_context_, impl_->memory_pool_,
impl_->executor_, impl_->specific_file_system_, impl_->table_schema_, impl_->options_,
impl_->cache_);
impl_->cache_, impl_->snapshot_read_view_);
impl_->Reset();
return ctx;
}
Expand Down
63 changes: 47 additions & 16 deletions src/paimon/core/table/source/data_evolution_batch_scan.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,7 @@
#include "paimon/core/global_index/global_index_scan_impl.h"
#include "paimon/core/global_index/indexed_split_impl.h"
#include "paimon/core/table/source/data_split_impl.h"
#include "paimon/core/table/source/snapshot_read_view_impl.h"
#include "paimon/core/utils/snapshot_manager.h"
#include "paimon/global_index/bitmap_global_index_result.h"

Expand All @@ -44,6 +45,11 @@ DataEvolutionBatchScan::DataEvolutionBatchScan(
executor_(executor) {}

Result<std::shared_ptr<Plan>> DataEvolutionBatchScan::CreatePlan() {
std::shared_ptr<const SnapshotReadView> snapshot_read_view =
snapshot_reader_->GetSnapshotReadView();
if (snapshot_read_view && !snapshot_read_view->SnapshotId()) {
return PlanImpl::EmptyPlan(snapshot_read_view);
}
std::optional<std::vector<Range>> row_ranges;
std::shared_ptr<GlobalIndexResult> final_global_index_result = global_index_result_;
if (!final_global_index_result) {
Expand All @@ -52,14 +58,22 @@ Result<std::shared_ptr<Plan>> DataEvolutionBatchScan::CreatePlan() {
final_global_index_result = index_result;
PAIMON_ASSIGN_OR_RAISE(row_ranges, index_result->ToRanges());
}
// EvalGlobalIndex captures the latest snapshot before consulting its index. Refresh the
// local copy so both negative and positive results publish and reuse that exact view.
snapshot_read_view = snapshot_reader_->GetSnapshotReadView();
} else {
PAIMON_ASSIGN_OR_RAISE(row_ranges, final_global_index_result->ToRanges());
}
if (!row_ranges) {
return batch_scan_->CreatePlan();
}
if (row_ranges.value().empty()) {
return PlanImpl::EmptyPlan();
if (snapshot_read_view && snapshot_read_view->SnapshotId()) {
return std::make_shared<PlanImpl>(snapshot_read_view->SnapshotId(),
std::vector<std::shared_ptr<Split>>(),
snapshot_read_view);
}
return PlanImpl::EmptyPlan(snapshot_read_view);
}
PAIMON_ASSIGN_OR_RAISE(RowRangeIndex row_range_index,
RowRangeIndex::Create(row_ranges.value()));
Expand Down Expand Up @@ -133,7 +147,8 @@ Result<std::shared_ptr<Plan>> DataEvolutionBatchScan::WrapToIndexedSplits(
}
indexed_splits.push_back(std::make_shared<IndexedSplitImpl>(data_split, expected, scores));
}
return std::make_shared<PlanImpl>(data_plan->SnapshotId(), indexed_splits);
return std::make_shared<PlanImpl>(data_plan->SnapshotId(), indexed_splits,
data_plan->GetSnapshotReadView());
}

Result<std::shared_ptr<GlobalIndexResult>> DataEvolutionBatchScan::EvalGlobalIndex() const {
Expand All @@ -145,24 +160,40 @@ Result<std::shared_ptr<GlobalIndexResult>> DataEvolutionBatchScan::EvalGlobalInd
return std::shared_ptr<GlobalIndexResult>(nullptr);
}
auto partition_filter = batch_scan_->GetPartitionPredicate();
// TODO(lisizhuo.lsz): support time travel
std::optional<Snapshot> snapshot;
const std::shared_ptr<SnapshotManager>& snapshot_manager =
snapshot_reader_->GetSnapshotManager();
if (const std::optional<int64_t>& snapshot_id = core_options_.GetScanSnapshotId()) {
PAIMON_ASSIGN_OR_RAISE(Snapshot loaded_snapshot,
snapshot_manager->LoadSnapshot(snapshot_id.value()));
snapshot = std::move(loaded_snapshot);
StartupMode startup_mode = core_options_.GetStartupMode();
std::shared_ptr<const Snapshot> snapshot;
if (!(startup_mode == StartupMode::LatestFull() || startup_mode == StartupMode::Latest())) {
Comment thread
wangyong9999 marked this conversation as resolved.
// Snapshot read views deliberately support latest scans only. Preserve main's non-latest
// planning path, including an explicitly selected snapshot and the caller-owned context.
// TODO(lisizhuo.lsz): support tag/timestamp time travel.
std::optional<Snapshot> loaded_snapshot;
const std::shared_ptr<SnapshotManager>& snapshot_manager =
snapshot_reader_->GetSnapshotManager();
if (const std::optional<int64_t>& snapshot_id = core_options_.GetScanSnapshotId()) {
PAIMON_ASSIGN_OR_RAISE(Snapshot selected_snapshot,
snapshot_manager->LoadSnapshot(snapshot_id.value()));
loaded_snapshot = std::move(selected_snapshot);
} else {
PAIMON_ASSIGN_OR_RAISE(loaded_snapshot, snapshot_manager->LatestSnapshot());
}
if (!loaded_snapshot) {
return Status::Invalid("not found latest snapshot");
}
snapshot = std::make_shared<const Snapshot>(std::move(loaded_snapshot).value());
} else {
PAIMON_ASSIGN_OR_RAISE(snapshot, snapshot_manager->LatestSnapshot());
}
if (!snapshot) {
return Status::Invalid("not found latest snapshot");
// Capture the latest snapshot through the shared SnapshotReader before evaluating the
// index. The subsequent data scan then consumes the same immutable view, avoiding both a
// second latest-snapshot lookup and a row-id/snapshot race.
PAIMON_ASSIGN_OR_RAISE(std::shared_ptr<const SnapshotReadView> snapshot_read_view,
snapshot_reader_->CaptureLatestSnapshotReadView());
PAIMON_ASSIGN_OR_RAISE(snapshot, SnapshotReadViewImpl::GetSnapshot(snapshot_read_view));
if (!snapshot) {
return std::shared_ptr<GlobalIndexResult>(nullptr);
}
}

PAIMON_ASSIGN_OR_RAISE(
std::unique_ptr<GlobalIndexScanImpl> index_scan,
GlobalIndexScanImpl::Create(table_path_, table_schema_, snapshot.value(), partition_filter,
GlobalIndexScanImpl::Create(table_path_, table_schema_, *snapshot, partition_filter,
core_options_, executor_, pool_));
return index_scan->Scan(predicate);
}
Expand Down
6 changes: 4 additions & 2 deletions src/paimon/core/table/source/data_table_batch_scan.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -90,7 +90,7 @@ Result<std::shared_ptr<Plan>> DataTableBatchScan::ApplyPushDownLimit(
std::dynamic_pointer_cast<StartingScanner::CurrentSnapshot>(scan_result);
if (!current_scan_result) {
// NoSnapshot
return PlanImpl::EmptyPlan();
return snapshot_reader_->EmptyPlan();
}
if (!CanPushDownLimit()) {
return current_scan_result->GetPlan();
Expand Down Expand Up @@ -126,7 +126,9 @@ Result<std::shared_ptr<Plan>> DataTableBatchScan::ApplyPushDownLimit(
"rows.",
limited_data_splits.size(), splits.size(), push_down_limit_.value(),
scanned_row_count);
return std::make_shared<PlanImpl>(snapshot_id, limited_data_splits);
return std::make_shared<PlanImpl>(
snapshot_id, limited_data_splits,
current_scan_result->GetPlan()->GetSnapshotReadView());
}
}
return current_scan_result->GetPlan();
Expand Down
9 changes: 9 additions & 0 deletions src/paimon/core/table/source/plan_impl.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -27,4 +27,13 @@ const std::shared_ptr<Plan> PlanImpl::EmptyPlan() {
return empty_plan;
}

const std::shared_ptr<Plan> PlanImpl::EmptyPlan(
const std::shared_ptr<const SnapshotReadView>& snapshot_read_view) {
if (!snapshot_read_view) {
return EmptyPlan();
}
return std::make_shared<PlanImpl>(std::optional<int64_t>(),
std::vector<std::shared_ptr<Split>>(), snapshot_read_view);
}

} // namespace paimon
15 changes: 14 additions & 1 deletion src/paimon/core/table/source/plan_impl.h
Original file line number Diff line number Diff line change
Expand Up @@ -33,7 +33,12 @@ class PlanImpl : public Plan {
public:
PlanImpl(const std::optional<int64_t>& snapshot_id,
const std::vector<std::shared_ptr<Split>>& splits)
: snapshot_id_(snapshot_id), splits_(splits) {}
: PlanImpl(snapshot_id, splits, nullptr) {}

PlanImpl(const std::optional<int64_t>& snapshot_id,
const std::vector<std::shared_ptr<Split>>& splits,
const std::shared_ptr<const SnapshotReadView>& snapshot_read_view)
: snapshot_id_(snapshot_id), splits_(splits), snapshot_read_view_(snapshot_read_view) {}

std::optional<int64_t> SnapshotId() const override {
return snapshot_id_;
Expand All @@ -43,10 +48,18 @@ class PlanImpl : public Plan {
return splits_;
}

std::shared_ptr<const SnapshotReadView> GetSnapshotReadView() const override {
return snapshot_read_view_;
}

static const std::shared_ptr<Plan> EmptyPlan();

static const std::shared_ptr<Plan> EmptyPlan(
const std::shared_ptr<const SnapshotReadView>& snapshot_read_view);

private:
std::optional<int64_t> snapshot_id_;
std::vector<std::shared_ptr<Split>> splits_;
std::shared_ptr<const SnapshotReadView> snapshot_read_view_;
};
} // namespace paimon
Loading