Skip to content
Merged
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
2 changes: 1 addition & 1 deletion contrib/tipb
21 changes: 16 additions & 5 deletions dbms/src/Debug/MockStorage.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -159,7 +159,7 @@ bool shouldEnableMultiStageLateMaterializationForMockDeltaMerge(
const TiDB::ColumnInfos & scan_column_infos)
{
const auto & settings = context.getSettingsRef();
if (!settings.dt_enable_multi_stage_late_materialization)
if (settings.dt_enable_multi_stage_late_materialization == 0)
return false;
if (filter_conditions == nullptr || !filter_conditions->hasValue())
return false;
Expand Down Expand Up @@ -376,7 +376,8 @@ BlockInputStreamPtr MockStorage::getStreamFromDeltaMerge(
std::vector<int> runtime_filter_ids,
int rf_max_wait_time_ms,
const google::protobuf::RepeatedPtrField<tipb::Expr> * pushed_down_filters,
const TiDB::ColumnInfos * table_scan_column_infos)
const TiDB::ColumnInfos * table_scan_column_infos,
const DM::MultiStageLateMaterializationTopNDescriptionPtr & multi_stage_late_materialization_topn)
{
static const google::protobuf::RepeatedPtrField<tipb::Expr> empty_pushed_down_filters{};
static const auto empty_ann_query_info = tipb::ANNQueryInfo{};
Expand All @@ -400,6 +401,7 @@ BlockInputStreamPtr MockStorage::getStreamFromDeltaMerge(
runtime_filter_ids,
rf_max_wait_time_ms,
context.getTimezoneInfo());
query_info.multi_stage_late_materialization_topn = multi_stage_late_materialization_topn;
BlockInputStreams ins = storage->read(
column_names,
query_info,
Expand Down Expand Up @@ -476,7 +478,8 @@ void MockStorage::buildExecFromDeltaMerge(
int rf_max_wait_time_ms,
const google::protobuf::RepeatedPtrField<tipb::Expr> * pushed_down_filters,
const String & table_scan_executor_id,
const TiDB::ColumnInfos * table_scan_column_infos)
const TiDB::ColumnInfos * table_scan_column_infos,
const DM::MultiStageLateMaterializationTopNDescriptionPtr & multi_stage_late_materialization_topn)
{
static const google::protobuf::RepeatedPtrField<tipb::Expr> empty_pushed_down_filters{};
static const auto empty_ann_query_info = tipb::ANNQueryInfo{};
Expand All @@ -501,9 +504,11 @@ void MockStorage::buildExecFromDeltaMerge(
if (enable_multi_stage_late_materialization)
{
multi_stage_late_materialization_runtime_stats
= std::make_shared<DM::MultiStageLateMaterializationRuntimeStats>();
= std::make_shared<DM::MultiStageLateMaterializationRuntimeStats>(
fmt::format("mock table_scan_executor_id={}", table_scan_executor_id));
if (auto * dag_context = context.getDAGContext(); dag_context != nullptr && !table_scan_executor_id.empty())
{
dag_context->scan_context_map[table_scan_executor_id] = query_info.mvcc_query_info->scan_context;
dag_context->setExecutorRowsOverride(
table_scan_executor_id,
std::shared_ptr<std::atomic<UInt64>>(
Expand All @@ -513,8 +518,10 @@ void MockStorage::buildExecFromDeltaMerge(
filter_conditions->executor_id,
std::shared_ptr<std::atomic<UInt64>>(
multi_stage_late_materialization_runtime_stats,
&multi_stage_late_materialization_runtime_stats->stage1_output_rows));
&multi_stage_late_materialization_runtime_stats->final_rest_input_rows));
}
query_info.mvcc_query_info->scan_context->setMultiStageLateMaterializationRuntimeStats(
multi_stage_late_materialization_runtime_stats);
}
query_info.dag_query = std::make_unique<DAGQueryInfo>(
filter_conditions->conditions,
Expand All @@ -526,6 +533,7 @@ void MockStorage::buildExecFromDeltaMerge(
context.getTimezoneInfo());
query_info.enable_multi_stage_late_materialization = enable_multi_stage_late_materialization;
query_info.multi_stage_late_materialization_runtime_stats = multi_stage_late_materialization_runtime_stats;
query_info.multi_stage_late_materialization_topn = multi_stage_late_materialization_topn;
storage->read(
exec_context_,
group_builder,
Expand Down Expand Up @@ -607,6 +615,8 @@ void MockStorage::addTableInfoForDeltaMerge(const String & name, const MockColum
TiDB::ColumnInfo ret;
ret.name = column.name;
ret.tp = column.type;
ret.collate = column.collate;
ret.elems = column.elems;

if (!column.nullable)
ret.setNotNullFlag();
Expand Down Expand Up @@ -890,6 +900,7 @@ TiDB::ColumnInfos mockColumnInfosToTiDBColumnInfos(const MockColumnInfoVec & moc
column_info.name = mock_column_info.name;
column_info.tp = mock_column_info.type;
column_info.collate = mock_column_info.collate;
column_info.elems = mock_column_info.elems;
column_info.id = col_id++;
// TODO: find a way to assign decimal field's flen.
if (column_info.tp == TiDB::TP::TypeNewDecimal)
Expand Down
25 changes: 23 additions & 2 deletions dbms/src/Debug/MockStorage.h
Original file line number Diff line number Diff line change
Expand Up @@ -20,12 +20,15 @@
#include <Flash/Pipeline/Exec/PipelineExecBuilder.h>
#include <Operators/Operator.h>
#include <Storages/DeltaMerge/ColumnDefine_fwd.h>
#include <Storages/DeltaMerge/MultiStageLateMaterializationTopN.h>
#include <TiDB/Schema/TiDB_fwd.h>
#include <common/types.h>

#include <atomic>
#include <memory>
#include <unordered_map>
#include <utility>
#include <vector>

namespace DB
{
Expand All @@ -36,10 +39,26 @@ struct SelectQueryInfo;

struct MockColumnInfo
{
MockColumnInfo() = default;

MockColumnInfo(
String name_,
TiDB::TP type_,
bool nullable_ = true,
Poco::Dynamic::Var collate_ = {},
std::vector<std::pair<std::string, Int16>> elems_ = {})
: name(std::move(name_))
, type(type_)
, nullable(nullable_)
, collate(std::move(collate_))
, elems(std::move(elems_))
{}

String name;
TiDB::TP type;
bool nullable = true;
Poco::Dynamic::Var collate{}; // default empty means no collation.
std::vector<std::pair<std::string, Int16>> elems;
};
using MockColumnInfoVec = std::vector<MockColumnInfo>;
using TableInfo = TiDB::TableInfo;
Expand Down Expand Up @@ -107,7 +126,8 @@ class MockStorage
std::vector<int> runtime_filter_ids = std::vector<int>(),
int rf_max_wait_time_ms = 0,
const google::protobuf::RepeatedPtrField<tipb::Expr> * pushed_down_filters = nullptr,
const TiDB::ColumnInfos * scan_column_infos = nullptr);
const TiDB::ColumnInfos * scan_column_infos = nullptr,
const DM::MultiStageLateMaterializationTopNDescriptionPtr & multi_stage_late_materialization_topn = nullptr);

void buildExecFromDeltaMerge(
PipelineExecutorContext & exec_context_,
Expand All @@ -121,7 +141,8 @@ class MockStorage
int rf_max_wait_time_ms = 0,
const google::protobuf::RepeatedPtrField<tipb::Expr> * pushed_down_filters = nullptr,
const String & table_scan_executor_id = "",
const TiDB::ColumnInfos * scan_column_infos = nullptr);
const TiDB::ColumnInfos * scan_column_infos = nullptr,
const DM::MultiStageLateMaterializationTopNDescriptionPtr & multi_stage_late_materialization_topn = nullptr);

bool tableExistsForDeltaMerge(Int64 table_id);

Expand Down
18 changes: 14 additions & 4 deletions dbms/src/Flash/Coprocessor/DAGStorageInterpreter.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -382,11 +382,13 @@ DAGStorageInterpreter::DAGStorageInterpreter(
Context & context_,
const TiDBTableScan & table_scan_,
const FilterConditions & filter_conditions_,
size_t max_streams_)
size_t max_streams_,
const DM::MultiStageLateMaterializationTopNDescriptionPtr & multi_stage_late_materialization_topn_)
: context(context_)
, table_scan(table_scan_)
, filter_conditions(filter_conditions_)
, max_streams(max_streams_)
, multi_stage_late_materialization_topn(multi_stage_late_materialization_topn_)
, log(Logger::get(context.getDAGContext()->log ? context.getDAGContext()->log->identifier() : ""))
, logical_table_id(table_scan.getLogicalTableID())
, tmt(context.getTMTContext())
Expand Down Expand Up @@ -1044,7 +1046,8 @@ std::unordered_map<TableID, SelectQueryInfo> DAGStorageInterpreter::generateSele
if (enable_multi_stage_late_materialization)
{
multi_stage_late_materialization_runtime_stats
= std::make_shared<DM::MultiStageLateMaterializationRuntimeStats>();
= std::make_shared<DM::MultiStageLateMaterializationRuntimeStats>(
fmt::format("{} table_scan_executor_id={}", log->identifier(), table_scan.getTableScanExecutorID()));
dagContext().setExecutorRowsOverride(
table_scan.getTableScanExecutorID(),
std::shared_ptr<std::atomic<UInt64>>(
Expand All @@ -1054,7 +1057,13 @@ std::unordered_map<TableID, SelectQueryInfo> DAGStorageInterpreter::generateSele
filter_conditions.executor_id,
std::shared_ptr<std::atomic<UInt64>>(
multi_stage_late_materialization_runtime_stats,
&multi_stage_late_materialization_runtime_stats->stage1_output_rows));
&multi_stage_late_materialization_runtime_stats->final_rest_input_rows));
if (auto scan_context_it = dagContext().scan_context_map.find(table_scan.getTableScanExecutorID());
scan_context_it != dagContext().scan_context_map.end() && scan_context_it->second != nullptr)
{
scan_context_it->second->setMultiStageLateMaterializationRuntimeStats(
multi_stage_late_materialization_runtime_stats);
}
}

auto create_query_info = [&](Int64 table_id) -> SelectQueryInfo {
Expand All @@ -1074,6 +1083,7 @@ std::unordered_map<TableID, SelectQueryInfo> DAGStorageInterpreter::generateSele
query_info.is_fast_scan = table_scan.isFastScan();
query_info.enable_multi_stage_late_materialization = enable_multi_stage_late_materialization;
query_info.multi_stage_late_materialization_runtime_stats = multi_stage_late_materialization_runtime_stats;
query_info.multi_stage_late_materialization_topn = multi_stage_late_materialization_topn;
return query_info;
};
RUNTIME_CHECK_MSG(mvcc_query_info->scan_context != nullptr, "Unexpected null scan_context");
Expand Down Expand Up @@ -1790,7 +1800,7 @@ std::pair<Names, std::vector<UInt8>> DAGStorageInterpreter::getColumnsForTableSc

bool DAGStorageInterpreter::shouldEnableMultiStageLateMaterialization() const
{
if (!context.getSettingsRef().dt_enable_multi_stage_late_materialization)
if (context.getSettingsRef().dt_enable_multi_stage_late_materialization == 0)
return false;

auto disable = [&](const String & reason) {
Expand Down
5 changes: 4 additions & 1 deletion dbms/src/Flash/Coprocessor/DAGStorageInterpreter.h
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,7 @@
#include <Flash/Coprocessor/RemoteRequest.h>
#include <Flash/Coprocessor/TiDBTableScan.h>
#include <Flash/Pipeline/Exec/PipelineExecBuilder.h>
#include <Storages/DeltaMerge/MultiStageLateMaterializationTopN.h>
#include <Storages/DeltaMerge/Remote/DisaggSnapshot_fwd.h>
#include <Storages/KVStore/Read/LearnerRead.h>
#include <Storages/KVStore/Read/RegionException.h>
Expand All @@ -47,7 +48,8 @@ class DAGStorageInterpreter
Context & context_,
const TiDBTableScan & table_scan,
const FilterConditions & filter_conditions_,
size_t max_streams_);
size_t max_streams_,
const DM::MultiStageLateMaterializationTopNDescriptionPtr & multi_stage_late_materialization_topn_ = nullptr);

~DAGStorageInterpreter();

Expand Down Expand Up @@ -146,6 +148,7 @@ class DAGStorageInterpreter
const TiDBTableScan & table_scan;
const FilterConditions & filter_conditions;
const size_t max_streams;
const DM::MultiStageLateMaterializationTopNDescriptionPtr multi_stage_late_materialization_topn;
LoggerPtr log;

/// derived from other members, doesn't change during DAGStorageInterpreter's lifetime
Expand Down
15 changes: 12 additions & 3 deletions dbms/src/Flash/Planner/Plans/PhysicalMockTableScan.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -73,7 +73,8 @@ std::pair<NamesAndTypes, BlockInputStreams> mockSchemaAndStreams(
table_scan.getRuntimeFilterIDs(),
10000,
&table_scan.getPushedDownFilters(),
use_table_scan_columns ? &table_scan.getColumns() : nullptr));
use_table_scan_columns ? &table_scan.getColumns() : nullptr,
nullptr));
}
else
{
Expand Down Expand Up @@ -184,7 +185,8 @@ void PhysicalMockTableScan::buildPipelineExecGroupImpl(
rf_max_wait_time_ms,
&pushed_down_filters,
execId(),
use_table_scan_columns_for_delta_merge ? &table_scan_columns : nullptr);
use_table_scan_columns_for_delta_merge ? &table_scan_columns : nullptr,
multi_stage_late_materialization_topn);
for (size_t i = 0; i < group_builder.concurrency(); ++i)
{
if (auto * source_op = dynamic_cast<UnorderedSourceOp *>(group_builder.getCurBuilder(i).source_op.get()))
Expand Down Expand Up @@ -240,7 +242,8 @@ bool PhysicalMockTableScan::setFilterConditions(
runtime_filter_ids,
10000,
&pushed_down_filters,
use_table_scan_columns_for_delta_merge ? &table_scan_columns : nullptr));
use_table_scan_columns_for_delta_merge ? &table_scan_columns : nullptr,
multi_stage_late_materialization_topn));

return true;
}
Expand All @@ -261,6 +264,12 @@ const String & PhysicalMockTableScan::getFilterConditionsId() const
return filter_conditions.executor_id;
}

void PhysicalMockTableScan::setMultiStageLateMaterializationTopN(
const DM::MultiStageLateMaterializationTopNDescriptionPtr & topn)
{
multi_stage_late_materialization_topn = topn;
}

void PhysicalMockTableScan::buildRuntimeFilterInLocalStream(Context & context)
{
for (const auto & local_stream : mock_streams)
Expand Down
9 changes: 9 additions & 0 deletions dbms/src/Flash/Planner/Plans/PhysicalMockTableScan.h
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@
#include <Flash/Coprocessor/RuntimeFilterMgr.h>
#include <Flash/Coprocessor/TiDBTableScan.h>
#include <Flash/Planner/Plans/PhysicalLeaf.h>
#include <Storages/DeltaMerge/MultiStageLateMaterializationTopN.h>
#include <tipb/executor.pb.h>

namespace DB
Expand Down Expand Up @@ -63,6 +64,12 @@ class PhysicalMockTableScan : public PhysicalLeaf

const String & getFilterConditionsId() const;

const TiDB::ColumnInfos & getTableScanColumns() const { return table_scan_columns; }

const google::protobuf::RepeatedPtrField<tipb::Expr> & getPushedDownFilters() const { return pushed_down_filters; }

void setMultiStageLateMaterializationTopN(const DM::MultiStageLateMaterializationTopNDescriptionPtr & topn);

private:
void buildBlockInputStreamImpl(DAGPipeline & pipeline, Context & /*context*/, size_t /*max_streams*/) override;

Expand Down Expand Up @@ -92,6 +99,8 @@ class PhysicalMockTableScan : public PhysicalLeaf

TiDB::ColumnInfos table_scan_columns;

DM::MultiStageLateMaterializationTopNDescriptionPtr multi_stage_late_materialization_topn;

bool use_table_scan_columns_for_delta_merge;

const int rf_max_wait_time_ms = 10000;
Expand Down
20 changes: 18 additions & 2 deletions dbms/src/Flash/Planner/Plans/PhysicalTableScan.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -114,7 +114,12 @@ void PhysicalTableScan::buildBlockInputStreamImpl(DAGPipeline & pipeline, Contex
}
else
{
DAGStorageInterpreter storage_interpreter(context, tidb_table_scan, filter_conditions, max_streams);
DAGStorageInterpreter storage_interpreter(
context,
tidb_table_scan,
filter_conditions,
max_streams,
multi_stage_late_materialization_topn);
storage_interpreter.execute(pipeline);
}
buildProjection(pipeline);
Expand All @@ -137,7 +142,12 @@ void PhysicalTableScan::buildPipeline(
}
else
{
DAGStorageInterpreter storage_interpreter(context, tidb_table_scan, filter_conditions, context.getMaxStreams());
DAGStorageInterpreter storage_interpreter(
context,
tidb_table_scan,
filter_conditions,
context.getMaxStreams(),
multi_stage_late_materialization_topn);
storage_interpreter.execute(exec_context, pipeline_exec_builder);
}
buildProjection(exec_context, pipeline_exec_builder);
Expand Down Expand Up @@ -214,4 +224,10 @@ const String & PhysicalTableScan::getFilterConditionsId() const
RUNTIME_CHECK(hasFilterConditions());
return filter_conditions.executor_id;
}

void PhysicalTableScan::setMultiStageLateMaterializationTopN(
const DM::MultiStageLateMaterializationTopNDescriptionPtr & topn)
{
multi_stage_late_materialization_topn = topn;
}
} // namespace DB
7 changes: 7 additions & 0 deletions dbms/src/Flash/Planner/Plans/PhysicalTableScan.h
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@
#include <Flash/Coprocessor/FilterConditions.h>
#include <Flash/Coprocessor/TiDBTableScan.h>
#include <Flash/Planner/Plans/PhysicalLeaf.h>
#include <Storages/DeltaMerge/MultiStageLateMaterializationTopN.h>
#include <tipb/executor.pb.h>

namespace DB
Expand Down Expand Up @@ -48,6 +49,10 @@ class PhysicalTableScan : public PhysicalLeaf

const String & getFilterConditionsId() const;

const TiDBTableScan & getTiDBTableScan() const { return tidb_table_scan; }

void setMultiStageLateMaterializationTopN(const DM::MultiStageLateMaterializationTopNDescriptionPtr & topn);

void buildPipeline(PipelineBuilder & builder, Context & context, PipelineExecutorContext & exec_context) override;

private:
Expand All @@ -67,6 +72,8 @@ class PhysicalTableScan : public PhysicalLeaf

TiDBTableScan tidb_table_scan;

DM::MultiStageLateMaterializationTopNDescriptionPtr multi_stage_late_materialization_topn;

Block sample_block;

PipelineExecGroupBuilder pipeline_exec_builder;
Expand Down
Loading