This is an automated email from the ASF dual-hosted git repository.
changchen pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/incubator-gluten.git
The following commit(s) were added to refs/heads/main by this push:
new 29375d17ad [GLUTEN-1632][CH]Daily Update Clickhouse Version (20241016)
(#7558)
29375d17ad is described below
commit 29375d17ad0c678b27b73635fa887b51bbc88083
Author: Kyligence Git <[email protected]>
AuthorDate: Wed Oct 16 09:12:36 2024 -0500
[GLUTEN-1632][CH]Daily Update Clickhouse Version (20241016) (#7558)
* [GLUTEN-1632][CH]Daily Update Clickhouse Version (20241016)
* Fix Build due to https://github.com/ClickHouse/ClickHouse/pull/70158
* Fix build due to https://github.com/ClickHouse/ClickHouse/pull/62966
* add more info for exception
* Fix ut due to https://github.com/ClickHouse/ClickHouse/pull/69481
---------
Co-authored-by: kyligence-git <[email protected]>
Co-authored-by: Chang Chen <[email protected]>
---
cpp-ch/clickhouse.version | 4 +-
cpp-ch/local-engine/Common/CHUtil.cpp | 21 +++---
.../Operator/BlocksBufferPoolTransform.cpp | 9 +--
.../Operator/BlocksBufferPoolTransform.h | 6 +-
.../Operator/DefaultHashAggregateResult.cpp | 7 +-
.../Operator/DefaultHashAggregateResult.h | 7 +-
cpp-ch/local-engine/Operator/EarlyStopStep.cpp | 8 +-
cpp-ch/local-engine/Operator/EarlyStopStep.h | 13 ++--
cpp-ch/local-engine/Operator/EmptyProjectStep.cpp | 10 +--
cpp-ch/local-engine/Operator/EmptyProjectStep.h | 6 +-
cpp-ch/local-engine/Operator/ExpandStep.cpp | 12 ++-
cpp-ch/local-engine/Operator/ExpandStep.h | 8 +-
.../local-engine/Operator/GraceAggregatingStep.cpp | 18 ++---
.../local-engine/Operator/GraceAggregatingStep.h | 13 ++--
.../Operator/GraceAggregatingTransform.cpp | 2 -
.../Operator/GraceAggregatingTransform.h | 5 +-
.../Operator/GraceMergingAggregatedStep.cpp | 28 +++----
.../Operator/GraceMergingAggregatedStep.h | 17 ++---
cpp-ch/local-engine/Operator/ReplicateRowsStep.cpp | 8 +-
cpp-ch/local-engine/Operator/ReplicateRowsStep.h | 15 ++--
.../Operator/StreamingAggregatingStep.cpp | 26 +++----
.../Operator/StreamingAggregatingStep.h | 18 ++---
.../local-engine/Operator/WindowGroupLimitStep.cpp | 12 +--
.../local-engine/Operator/WindowGroupLimitStep.h | 8 +-
cpp-ch/local-engine/Parser/InputFileNameParser.cpp | 23 +++---
cpp-ch/local-engine/Parser/InputFileNameParser.h | 2 +-
cpp-ch/local-engine/Parser/LocalExecutor.cpp | 2 +-
.../Parser/RelParsers/AggregateRelParser.cpp | 42 +++++------
.../Parser/RelParsers/CrossRelParser.cpp | 54 +++++++-------
.../Parser/RelParsers/ExpandRelParser.cpp | 4 +-
.../Parser/RelParsers/FetchRelParser.cpp | 2 +-
.../Parser/RelParsers/FilterRelParser.cpp | 8 +-
.../Parser/RelParsers/JoinRelParser.cpp | 87 ++++++++++------------
.../Parser/RelParsers/MergeTreeRelParser.cpp | 2 +-
.../Parser/RelParsers/ProjectRelParser.cpp | 22 +++---
.../Parser/RelParsers/ReadRelParser.cpp | 2 +-
.../Parser/RelParsers/SortRelParser.cpp | 10 +--
.../RelParsers/WindowGroupLimitRelParser.cpp | 18 +----
.../Parser/RelParsers/WindowRelParser.cpp | 20 ++---
.../local-engine/Parser/SerializedPlanParser.cpp | 20 +++--
.../local-engine/Parser/SubstraitParserUtils.cpp | 16 +++-
cpp-ch/local-engine/Parser/SubstraitParserUtils.h | 2 +
.../Parquet/VectorizedParquetRecordReader.cpp | 14 +++-
.../SubstraitSource/SubstraitFileSourceStep.cpp | 4 +-
.../local-engine/tests/benchmark_local_engine.cpp | 22 +++---
cpp-ch/local-engine/tests/gtest_ch_join.cpp | 16 ++--
46 files changed, 321 insertions(+), 352 deletions(-)
diff --git a/cpp-ch/clickhouse.version b/cpp-ch/clickhouse.version
index 8c25141d6d..63308f5a51 100644
--- a/cpp-ch/clickhouse.version
+++ b/cpp-ch/clickhouse.version
@@ -1,3 +1,3 @@
CH_ORG=Kyligence
-CH_BRANCH=rebase_ch/20241015
-CH_COMMIT=7e3b2d69a74
+CH_BRANCH=rebase_ch/20241016
+CH_COMMIT=c3a1c1b8457
diff --git a/cpp-ch/local-engine/Common/CHUtil.cpp
b/cpp-ch/local-engine/Common/CHUtil.cpp
index 2615adf08a..44d86168d6 100644
--- a/cpp-ch/local-engine/Common/CHUtil.cpp
+++ b/cpp-ch/local-engine/Common/CHUtil.cpp
@@ -49,7 +49,6 @@
#include <IO/SharedThreadPools.h>
#include <Interpreters/JIT/CompiledExpressionCache.h>
#include <Parser/RelParsers/RelParser.h>
-#include <Parser/SubstraitParserUtils.h>
#include <Planner/PlannerActionsVisitor.h>
#include <Processors/Chunk.h>
#include <Processors/QueryPlan/ExpressionStep.h>
@@ -354,11 +353,10 @@ void PlanUtil::checkOuputType(const DB::QueryPlan & plan)
// QueryPlan::checkInitialized is a private method, so we assume plan is
initialized, otherwise there is a core dump here.
// It's okay, because it's impossible for us not to initialize where we
call this method.
const auto & step = *plan.getRootNode()->step;
- if (!step.hasOutputStream())
- return;
- if (!step.getOutputStream().header)
+
+ if (!step.hasOutputHeader())
return;
- for (const auto & elem : step.getOutputStream().header)
+ for (const auto & elem : step.getOutputHeader())
{
const DB::DataTypePtr & ch_type = elem.type;
const auto ch_type_without_nullable = DB::removeNullable(ch_type);
@@ -385,10 +383,10 @@ void PlanUtil::checkOuputType(const DB::QueryPlan & plan)
DB::IQueryPlanStep * PlanUtil::adjustQueryPlanHeader(DB::QueryPlan & plan,
const DB::Block & to_header, const String & step_desc)
{
auto convert_actions_dag = DB::ActionsDAG::makeConvertingActions(
- plan.getCurrentDataStream().header.getColumnsWithTypeAndName(),
+ plan.getCurrentHeader().getColumnsWithTypeAndName(),
to_header.getColumnsWithTypeAndName(),
ActionsDAG::MatchColumnsMode::Name);
- auto expression_step =
std::make_unique<DB::ExpressionStep>(plan.getCurrentDataStream(),
std::move(convert_actions_dag));
+ auto expression_step =
std::make_unique<DB::ExpressionStep>(plan.getCurrentHeader(),
std::move(convert_actions_dag));
expression_step->setStepDescription(step_desc);
auto * step_ptr = expression_step.get();
plan.addStep(std::move(expression_step));
@@ -399,7 +397,7 @@ DB::IQueryPlanStep *
PlanUtil::addRemoveNullableStep(DB::ContextPtr context, DB:
{
if (columns.empty())
return nullptr;
- DB::ActionsDAG
remove_nullable_actions_dag{plan.getCurrentDataStream().header.getColumnsWithTypeAndName()};
+ DB::ActionsDAG
remove_nullable_actions_dag{plan.getCurrentHeader().getColumnsWithTypeAndName()};
for (const auto & col_name : columns)
{
if (const auto * required_node =
remove_nullable_actions_dag.tryFindInOutputs(col_name))
@@ -410,7 +408,7 @@ DB::IQueryPlanStep *
PlanUtil::addRemoveNullableStep(DB::ContextPtr context, DB:
remove_nullable_actions_dag.addOrReplaceInOutputs(node);
}
}
- auto expression_step =
std::make_unique<DB::ExpressionStep>(plan.getCurrentDataStream(),
std::move(remove_nullable_actions_dag));
+ auto expression_step =
std::make_unique<DB::ExpressionStep>(plan.getCurrentHeader(),
std::move(remove_nullable_actions_dag));
expression_step->setStepDescription("Remove nullable properties");
auto * step_ptr = expression_step.get();
plan.addStep(std::move(expression_step));
@@ -783,6 +781,7 @@ void BackendInitializerUtil::initSettings(const
std::map<std::string, std::strin
settings.set("precise_float_parsing", true);
settings.set("enable_named_columns_in_function_tuple", false);
settings.set("date_time_64_output_format_cut_trailing_zeros_align_to_groups_of_thousands",
true);
+ settings.set("input_format_orc_dictionary_as_low_cardinality", false);
//after https://github.com/ClickHouse/ClickHouse/pull/69481
if (spark_conf_map.contains(GLUTEN_TASK_OFFHEAP))
{
@@ -1045,14 +1044,14 @@ UInt64 MemoryUtil::getMemoryRSS()
void JoinUtil::reorderJoinOutput(DB::QueryPlan & plan, DB::Names cols)
{
- ActionsDAG
project{plan.getCurrentDataStream().header.getNamesAndTypesList()};
+ ActionsDAG project{plan.getCurrentHeader().getNamesAndTypesList()};
NamesWithAliases project_cols;
for (const auto & col : cols)
{
project_cols.emplace_back(NameWithAlias(col, col));
}
project.project(project_cols);
- QueryPlanStepPtr project_step =
std::make_unique<ExpressionStep>(plan.getCurrentDataStream(),
std::move(project));
+ QueryPlanStepPtr project_step =
std::make_unique<ExpressionStep>(plan.getCurrentHeader(), std::move(project));
project_step->setStepDescription("Reorder Join Output");
plan.addStep(std::move(project_step));
}
diff --git a/cpp-ch/local-engine/Operator/BlocksBufferPoolTransform.cpp
b/cpp-ch/local-engine/Operator/BlocksBufferPoolTransform.cpp
index 16a5bd5d26..fa3dbf6431 100644
--- a/cpp-ch/local-engine/Operator/BlocksBufferPoolTransform.cpp
+++ b/cpp-ch/local-engine/Operator/BlocksBufferPoolTransform.cpp
@@ -84,9 +84,8 @@ void BlocksBufferPoolTransform::work()
{
}
-BlocksBufferPoolStep::BlocksBufferPoolStep(const DB::DataStream &
input_stream_, size_t buffer_size_)
- : DB::ITransformingStep(input_stream_, input_stream_.header, getTraits())
- , header(input_stream_.header)
+BlocksBufferPoolStep::BlocksBufferPoolStep(const DB::Block & input_header,
size_t buffer_size_)
+ : DB::ITransformingStep(input_header, input_header, getTraits())
, buffer_size(buffer_size_)
{
}
@@ -112,9 +111,9 @@ void
BlocksBufferPoolStep::describePipeline(DB::IQueryPlanStep::FormatSettings &
DB::IQueryPlanStep::describePipeline(processors, settings);
}
-void BlocksBufferPoolStep::updateOutputStream()
+void BlocksBufferPoolStep::updateOutputHeader()
{
- createOutputStream(input_streams.front(), input_streams.front().header,
getDataStreamTraits());
+ output_header = input_headers.front();
}
}
diff --git a/cpp-ch/local-engine/Operator/BlocksBufferPoolTransform.h
b/cpp-ch/local-engine/Operator/BlocksBufferPoolTransform.h
index 0ddc8a7403..787052d093 100644
--- a/cpp-ch/local-engine/Operator/BlocksBufferPoolTransform.h
+++ b/cpp-ch/local-engine/Operator/BlocksBufferPoolTransform.h
@@ -22,13 +22,12 @@
#include <Processors/QueryPlan/IQueryPlanStep.h>
#include <Processors/QueryPlan/ITransformingStep.h>
-
namespace local_engine
{
class BlocksBufferPoolStep : public DB::ITransformingStep
{
public:
- explicit BlocksBufferPoolStep(const DB::DataStream & input_stream_, size_t
buffer_size_ = 4);
+ explicit BlocksBufferPoolStep(const DB::Block & input_header, size_t
buffer_size_ = 4);
~BlocksBufferPoolStep() override = default;
String getName() const override { return "BlocksBufferPoolStep"; }
@@ -36,9 +35,8 @@ public:
void transformPipeline(DB::QueryPipelineBuilder & pipeline, const
DB::BuildQueryPipelineSettings & settings) override;
void describePipeline(DB::IQueryPlanStep::FormatSettings & settings) const
override;
private:
- DB::Block header;
size_t buffer_size;
- void updateOutputStream() override;
+ void updateOutputHeader() override;
};
class BlocksBufferPoolTransform : public DB::IProcessor
diff --git a/cpp-ch/local-engine/Operator/DefaultHashAggregateResult.cpp
b/cpp-ch/local-engine/Operator/DefaultHashAggregateResult.cpp
index e96dab28bd..9278a7f70c 100644
--- a/cpp-ch/local-engine/Operator/DefaultHashAggregateResult.cpp
+++ b/cpp-ch/local-engine/Operator/DefaultHashAggregateResult.cpp
@@ -139,8 +139,8 @@ private:
DB::Chunk output_chunk;
};
-DefaultHashAggregateResultStep::DefaultHashAggregateResultStep(const
DB::DataStream & input_stream_)
- : DB::ITransformingStep(input_stream_,
adjustOutputHeader(input_stream_.header), getTraits())
+DefaultHashAggregateResultStep::DefaultHashAggregateResultStep(const DB::Block
& input_header)
+ : DB::ITransformingStep(input_header, adjustOutputHeader(input_header),
getTraits())
{
}
@@ -169,8 +169,7 @@ void
DefaultHashAggregateResultStep::describePipeline(DB::IQueryPlanStep::Format
DB::IQueryPlanStep::describePipeline(processors, settings);
}
-void DefaultHashAggregateResultStep::updateOutputStream()
+void DefaultHashAggregateResultStep::updateOutputHeader()
{
- createOutputStream(input_streams.front(), input_streams.front().header,
getDataStreamTraits());
}
}
diff --git a/cpp-ch/local-engine/Operator/DefaultHashAggregateResult.h
b/cpp-ch/local-engine/Operator/DefaultHashAggregateResult.h
index b433f40729..645ee2356c 100644
--- a/cpp-ch/local-engine/Operator/DefaultHashAggregateResult.h
+++ b/cpp-ch/local-engine/Operator/DefaultHashAggregateResult.h
@@ -16,7 +16,6 @@
*/
#pragma once
#include <Core/Block.h>
-#include <Parser/ExpandField.h>
#include <Processors/QueryPlan/IQueryPlanStep.h>
#include <Processors/QueryPlan/ITransformingStep.h>
namespace local_engine
@@ -26,13 +25,15 @@ namespace local_engine
class DefaultHashAggregateResultStep : public DB::ITransformingStep
{
public:
- explicit DefaultHashAggregateResultStep(const DB::DataStream &
input_stream_);
+ explicit DefaultHashAggregateResultStep(const DB::Block & input_header);
~DefaultHashAggregateResultStep() override = default;
String getName() const override { return "DefaultHashAggregateResultStep";
}
void transformPipeline(DB::QueryPipelineBuilder & pipeline, const
DB::BuildQueryPipelineSettings & settings) override;
void describePipeline(DB::IQueryPlanStep::FormatSettings & settings) const
override;
- void updateOutputStream() override;
+
+protected:
+ void updateOutputHeader() override;
};
}
diff --git a/cpp-ch/local-engine/Operator/EarlyStopStep.cpp
b/cpp-ch/local-engine/Operator/EarlyStopStep.cpp
index ff148a203b..d50ebd45f6 100644
--- a/cpp-ch/local-engine/Operator/EarlyStopStep.cpp
+++ b/cpp-ch/local-engine/Operator/EarlyStopStep.cpp
@@ -42,9 +42,9 @@ static DB::ITransformingStep::Traits getTraits()
}
EarlyStopStep::EarlyStopStep(
- const DB::DataStream & input_stream_)
+ const DB::Block & input_header_)
: DB::ITransformingStep(
- input_stream_, transformHeader(input_stream_.header), getTraits())
+ input_header_, transformHeader(input_header_), getTraits())
{
}
@@ -68,9 +68,9 @@ void
EarlyStopStep::describeActions(DB::IQueryPlanStep::FormatSettings & setting
DB::IQueryPlanStep::describePipeline(processors, settings);
}
-void EarlyStopStep::updateOutputStream()
+void EarlyStopStep::updateOutputHeader()
{
- output_stream = createOutputStream(input_streams.front(),
transformHeader(input_streams.front().header), getDataStreamTraits());
+ output_header = transformHeader(input_headers.front());
}
EarlyStopTransform::EarlyStopTransform(const DB::Block &header_)
diff --git a/cpp-ch/local-engine/Operator/EarlyStopStep.h
b/cpp-ch/local-engine/Operator/EarlyStopStep.h
index becfd46b83..bac8f67c64 100644
--- a/cpp-ch/local-engine/Operator/EarlyStopStep.h
+++ b/cpp-ch/local-engine/Operator/EarlyStopStep.h
@@ -14,6 +14,8 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
+
+#pragma once
#include <Core/Block.h>
#include <Processors/IProcessor.h>
#include <Processors/QueryPlan/IQueryPlanStep.h>
@@ -25,31 +27,32 @@ namespace local_engine
class EarlyStopStep : public DB::ITransformingStep
{
public:
- explicit EarlyStopStep(
- const DB::DataStream & input_stream_);
+ explicit EarlyStopStep(const DB::Block & input_header_);
~EarlyStopStep() override = default;
String getName() const override { return "EarlyStopStep"; }
- static DB::Block transformHeader(const DB::Block& input);
+ static DB::Block transformHeader(const DB::Block & input);
void transformPipeline(DB::QueryPipelineBuilder & pipeline, const
DB::BuildQueryPipelineSettings &) override;
void describeActions(DB::IQueryPlanStep::FormatSettings & settings) const
override;
+
private:
- void updateOutputStream() override;
+ void updateOutputHeader() override;
};
class EarlyStopTransform : public DB::IProcessor
{
public:
using Status = DB::IProcessor::Status;
- explicit EarlyStopTransform(const DB::Block &header_);
+ explicit EarlyStopTransform(const DB::Block & header_);
~EarlyStopTransform() override = default;
Status prepare() override;
void work() override;
String getName() const override { return "EarlyStopTransform"; }
+
private:
DB::Block header;
};
diff --git a/cpp-ch/local-engine/Operator/EmptyProjectStep.cpp
b/cpp-ch/local-engine/Operator/EmptyProjectStep.cpp
index 62991585f3..7c45556118 100644
--- a/cpp-ch/local-engine/Operator/EmptyProjectStep.cpp
+++ b/cpp-ch/local-engine/Operator/EmptyProjectStep.cpp
@@ -15,11 +15,11 @@
* limitations under the License.
*/
#include "EmptyProjectStep.h"
-#include <Common/CHUtil.h>
#include <Processors/Chunk.h>
#include <Processors/IProcessor.h>
#include <QueryPipeline/Pipe.h>
#include <QueryPipeline/QueryPipelineBuilder.h>
+#include <Common/CHUtil.h>
namespace local_engine
{
@@ -81,6 +81,7 @@ public:
has_input = false;
has_output = true;
}
+
private:
DB::Chunk output_chunk;
bool has_input = false;
@@ -100,8 +101,8 @@ static DB::ITransformingStep::Traits getTraits()
}};
}
-EmptyProjectStep::EmptyProjectStep(const DB::DataStream & input_stream_)
- : ITransformingStep(input_stream_, BlockUtil::buildRowCountHeader(),
getTraits())
+EmptyProjectStep::EmptyProjectStep(const DB::Block & input_header)
+ : ITransformingStep(input_header, BlockUtil::buildRowCountHeader(),
getTraits())
{
}
@@ -127,8 +128,7 @@ void
EmptyProjectStep::describePipeline(DB::IQueryPlanStep::FormatSettings & set
DB::IQueryPlanStep::describePipeline(processors, settings);
}
-void EmptyProjectStep::updateOutputStream()
+void EmptyProjectStep::updateOutputHeader()
{
- createOutputStream(input_streams.front(),
BlockUtil::buildRowCountHeader(), getDataStreamTraits());
}
}
diff --git a/cpp-ch/local-engine/Operator/EmptyProjectStep.h
b/cpp-ch/local-engine/Operator/EmptyProjectStep.h
index b2b42006f1..d1be3f3f93 100644
--- a/cpp-ch/local-engine/Operator/EmptyProjectStep.h
+++ b/cpp-ch/local-engine/Operator/EmptyProjectStep.h
@@ -28,13 +28,15 @@ namespace local_engine
class EmptyProjectStep : public DB::ITransformingStep
{
public:
- explicit EmptyProjectStep(const DB::DataStream & input_stream_);
+ explicit EmptyProjectStep(const DB::Block & input_header);
~EmptyProjectStep() override = default;
String getName() const override { return "EmptyProjectStep"; }
void transformPipeline(DB::QueryPipelineBuilder & pipeline, const
DB::BuildQueryPipelineSettings & settings) override;
void describePipeline(DB::IQueryPlanStep::FormatSettings & settings) const
override;
- void updateOutputStream() override;
+
+protected:
+ void updateOutputHeader() override;
};
}
diff --git a/cpp-ch/local-engine/Operator/ExpandStep.cpp
b/cpp-ch/local-engine/Operator/ExpandStep.cpp
index 518f094cdb..c9eccd2190 100644
--- a/cpp-ch/local-engine/Operator/ExpandStep.cpp
+++ b/cpp-ch/local-engine/Operator/ExpandStep.cpp
@@ -44,12 +44,10 @@ static DB::ITransformingStep::Traits getTraits()
}};
}
-ExpandStep::ExpandStep(const DB::DataStream & input_stream_, const ExpandField
& project_set_exprs_)
- : DB::ITransformingStep(input_stream_,
buildOutputHeader(input_stream_.header, project_set_exprs_), getTraits())
+ExpandStep::ExpandStep(const DB::Block & input_header, const ExpandField &
project_set_exprs_)
+ : DB::ITransformingStep(input_header, buildOutputHeader(input_header,
project_set_exprs_), getTraits())
, project_set_exprs(project_set_exprs_)
{
- header = input_stream_.header;
- output_header = getOutputStream().header;
}
DB::Block ExpandStep::buildOutputHeader(const DB::Block &, const ExpandField &
project_set_exprs_)
@@ -73,7 +71,7 @@ void ExpandStep::transformPipeline(DB::QueryPipelineBuilder &
pipeline, const DB
DB::Processors new_processors;
for (auto & output : outputs)
{
- auto expand_op = std::make_shared<ExpandTransform>(header,
output_header, project_set_exprs);
+ auto expand_op =
std::make_shared<ExpandTransform>(input_headers.front(), *output_header,
project_set_exprs);
new_processors.push_back(expand_op);
DB::connect(*output, expand_op->getInputs().front());
}
@@ -88,9 +86,9 @@ void
ExpandStep::describePipeline(DB::IQueryPlanStep::FormatSettings & settings)
DB::IQueryPlanStep::describePipeline(processors, settings);
}
-void ExpandStep::updateOutputStream()
+void ExpandStep::updateOutputHeader()
{
- createOutputStream(input_streams.front(), output_header,
getDataStreamTraits());
+ output_header = buildOutputHeader(input_headers.front(),
project_set_exprs);
}
}
diff --git a/cpp-ch/local-engine/Operator/ExpandStep.h
b/cpp-ch/local-engine/Operator/ExpandStep.h
index d35da5c4e9..810446674e 100644
--- a/cpp-ch/local-engine/Operator/ExpandStep.h
+++ b/cpp-ch/local-engine/Operator/ExpandStep.h
@@ -27,7 +27,7 @@ class ExpandStep : public DB::ITransformingStep
{
public:
// The input stream should only contain grouping columns.
- explicit ExpandStep(const DB::DataStream & input_stream_, const
ExpandField & project_set_exprs_);
+ explicit ExpandStep(const DB::Block & input_header, const ExpandField &
project_set_exprs_);
~ExpandStep() override = default;
String getName() const override { return "ExpandStep"; }
@@ -35,12 +35,10 @@ public:
void transformPipeline(DB::QueryPipelineBuilder & pipeline, const
DB::BuildQueryPipelineSettings & settings) override;
void describePipeline(DB::IQueryPlanStep::FormatSettings & settings) const
override;
-private:
+protected:
ExpandField project_set_exprs;
- DB::Block header;
- DB::Block output_header;
- void updateOutputStream() override;
+ void updateOutputHeader() override;
static DB::Block buildOutputHeader(const DB::Block & header, const
ExpandField & project_set_exprs_);
};
diff --git a/cpp-ch/local-engine/Operator/GraceAggregatingStep.cpp
b/cpp-ch/local-engine/Operator/GraceAggregatingStep.cpp
index 63175657f4..83c74bdd28 100644
--- a/cpp-ch/local-engine/Operator/GraceAggregatingStep.cpp
+++ b/cpp-ch/local-engine/Operator/GraceAggregatingStep.cpp
@@ -17,15 +17,10 @@
#include "GraceAggregatingStep.h"
#include <Interpreters/JoinUtils.h>
+#include <Operator/GraceAggregatingTransform.h>
#include <Processors/Transforms/AggregatingTransform.h>
#include <QueryPipeline/QueryPipelineBuilder.h>
-#include <Common/BitHelpers.h>
#include <Common/CHUtil.h>
-#include <Common/CurrentThread.h>
-#include <Common/GlutenConfig.h>
-#include <Common/QueryContext.h>
-#include <Common/formatReadable.h>
-#include "GraceAggregatingTransform.h"
namespace DB::ErrorCodes
{
@@ -52,8 +47,8 @@ static DB::Block buildOutputHeader(const DB::Block &
input_header_, const DB::Ag
}
GraceAggregatingStep::GraceAggregatingStep(
- DB::ContextPtr context_, const DB::DataStream & input_stream_,
DB::Aggregator::Params params_, bool no_pre_aggregated_)
- : DB::ITransformingStep(input_stream_,
buildOutputHeader(input_stream_.header, params_, false), getTraits())
+ DB::ContextPtr context_, const DB::Block & input_header,
DB::Aggregator::Params params_, bool no_pre_aggregated_)
+ : DB::ITransformingStep(input_header, buildOutputHeader(input_header,
params_, false), getTraits())
, context(context_)
, params(std::move(params_))
, no_pre_aggregated(no_pre_aggregated_)
@@ -87,7 +82,7 @@ void
GraceAggregatingStep::transformPipeline(DB::QueryPipelineBuilder & pipeline
void GraceAggregatingStep::describeActions(DB::IQueryPlanStep::FormatSettings
& settings) const
{
- return params.explain(settings.out, settings.offset);
+ params.explain(settings.out, settings.offset);
}
void GraceAggregatingStep::describeActions(DB::JSONBuilder::JSONMap & map)
const
@@ -95,10 +90,9 @@ void
GraceAggregatingStep::describeActions(DB::JSONBuilder::JSONMap & map) const
params.explain(map);
}
-void GraceAggregatingStep::updateOutputStream()
+void GraceAggregatingStep::updateOutputHeader()
{
- output_stream
- = createOutputStream(input_streams.front(),
buildOutputHeader(input_streams.front().header, params, false),
getDataStreamTraits());
+ output_header = buildOutputHeader(input_headers.front(), params, false);
}
diff --git a/cpp-ch/local-engine/Operator/GraceAggregatingStep.h
b/cpp-ch/local-engine/Operator/GraceAggregatingStep.h
index 7bb3917c9f..3ee069a02a 100644
--- a/cpp-ch/local-engine/Operator/GraceAggregatingStep.h
+++ b/cpp-ch/local-engine/Operator/GraceAggregatingStep.h
@@ -14,11 +14,10 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
+#pragma once
#include <Interpreters/Aggregator.h>
#include <Interpreters/Context.h>
-#include <Processors/IProcessor.h>
-#include <Processors/Port.h>
#include <Processors/QueryPlan/IQueryPlanStep.h>
#include <Processors/QueryPlan/ITransformingStep.h>
#include <Processors/Transforms/AggregatingTransform.h>
@@ -33,10 +32,7 @@ class GraceAggregatingStep : public DB::ITransformingStep
{
public:
explicit GraceAggregatingStep(
- DB::ContextPtr context_,
- const DB::DataStream & input_stream_,
- DB::Aggregator::Params params_,
- bool no_pre_aggregated_);
+ DB::ContextPtr context_, const DB::Block & input_header,
DB::Aggregator::Params params_, bool no_pre_aggregated_);
~GraceAggregatingStep() override = default;
String getName() const override { return "GraceAggregatingStep"; }
@@ -45,10 +41,13 @@ public:
void describeActions(DB::JSONBuilder::JSONMap & map) const override;
void describeActions(DB::IQueryPlanStep::FormatSettings & settings) const
override;
+
private:
DB::ContextPtr context;
DB::Aggregator::Params params;
bool no_pre_aggregated;
- void updateOutputStream() override;
+
+protected:
+ void updateOutputHeader() override;
};
}
\ No newline at end of file
diff --git a/cpp-ch/local-engine/Operator/GraceAggregatingTransform.cpp
b/cpp-ch/local-engine/Operator/GraceAggregatingTransform.cpp
index c8797d48d9..adf25d13f2 100644
--- a/cpp-ch/local-engine/Operator/GraceAggregatingTransform.cpp
+++ b/cpp-ch/local-engine/Operator/GraceAggregatingTransform.cpp
@@ -23,8 +23,6 @@
#include <Common/QueryContext.h>
#include <Common/formatReadable.h>
-#include <Common/DebugUtils.h>
-
namespace DB::ErrorCodes
{
extern const int LOGICAL_ERROR;
diff --git a/cpp-ch/local-engine/Operator/GraceAggregatingTransform.h
b/cpp-ch/local-engine/Operator/GraceAggregatingTransform.h
index 50006ea7d7..c2b787393a 100644
--- a/cpp-ch/local-engine/Operator/GraceAggregatingTransform.h
+++ b/cpp-ch/local-engine/Operator/GraceAggregatingTransform.h
@@ -14,18 +14,17 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
+#pragma once
#include <Core/Block.h>
#include <Interpreters/Aggregator.h>
#include <Interpreters/Context.h>
#include <Interpreters/JoinUtils.h>
#include <Processors/Chunk.h>
#include <Processors/IProcessor.h>
-#include <Processors/Port.h>
-#include <Processors/QueryPlan/ITransformingStep.h>
#include <Processors/Transforms/AggregatingTransform.h>
#include <Poco/Logger.h>
#include <Common/AggregateUtil.h>
-#include <Common/logger_useful.h>
+
namespace local_engine
{
diff --git a/cpp-ch/local-engine/Operator/GraceMergingAggregatedStep.cpp
b/cpp-ch/local-engine/Operator/GraceMergingAggregatedStep.cpp
index 348629f8aa..d6010f5b6f 100644
--- a/cpp-ch/local-engine/Operator/GraceMergingAggregatedStep.cpp
+++ b/cpp-ch/local-engine/Operator/GraceMergingAggregatedStep.cpp
@@ -15,14 +15,10 @@
* limitations under the License.
*/
#include "GraceMergingAggregatedStep.h"
-#include <Interpreters/JoinUtils.h>
+#include <Operator/GraceAggregatingTransform.h>
#include <Processors/Transforms/AggregatingTransform.h>
#include <QueryPipeline/QueryPipelineBuilder.h>
#include <Common/CHUtil.h>
-#include <Common/CurrentThread.h>
-#include <Common/formatReadable.h>
-#include <Common/BitHelpers.h>
-#include <Common/GlutenConfig.h>
#include <Common/QueryContext.h>
namespace DB
@@ -37,16 +33,14 @@ namespace local_engine
{
static DB::ITransformingStep::Traits getTraits()
{
- return DB::ITransformingStep::Traits
- {
+ return DB::ITransformingStep::Traits{
{
.preserves_number_of_streams = false,
.preserves_sorting = false,
},
{
.preserves_number_of_rows = false,
- }
- };
+ }};
}
static DB::Block buildOutputHeader(const DB::Block & input_header_, const
DB::Aggregator::Params params_, bool final)
@@ -55,12 +49,8 @@ static DB::Block buildOutputHeader(const DB::Block &
input_header_, const DB::Ag
}
GraceMergingAggregatedStep::GraceMergingAggregatedStep(
- DB::ContextPtr context_,
- const DB::DataStream & input_stream_,
- DB::Aggregator::Params params_,
- bool no_pre_aggregated_)
- : DB::ITransformingStep(
- input_stream_, buildOutputHeader(input_stream_.header, params_, true),
getTraits())
+ DB::ContextPtr context_, const DB::Block & input_header,
DB::Aggregator::Params params_, bool no_pre_aggregated_)
+ : DB::ITransformingStep(input_header, buildOutputHeader(input_header,
params_, true), getTraits())
, context(context_)
, params(std::move(params_))
, no_pre_aggregated(no_pre_aggregated_)
@@ -71,7 +61,8 @@ void
GraceMergingAggregatedStep::transformPipeline(DB::QueryPipelineBuilder & pi
{
if (params.max_bytes_before_external_group_by)
{
- throw DB::Exception(DB::ErrorCodes::LOGICAL_ERROR,
"max_bytes_before_external_group_by is not supported in
GraceMergingAggregatedStep");
+ throw DB::Exception(
+ DB::ErrorCodes::LOGICAL_ERROR, "max_bytes_before_external_group_by
is not supported in GraceMergingAggregatedStep");
}
auto num_streams = pipeline.getNumStreams();
auto transform_params =
std::make_shared<DB::AggregatingTransformParams>(pipeline.getHeader(), params,
true);
@@ -101,10 +92,9 @@ void
GraceMergingAggregatedStep::describeActions(DB::JSONBuilder::JSONMap & map)
params.explain(map);
}
-void GraceMergingAggregatedStep::updateOutputStream()
+void GraceMergingAggregatedStep::updateOutputHeader()
{
- output_stream = createOutputStream(input_streams.front(),
buildOutputHeader(input_streams.front().header, params, true),
getDataStreamTraits());
+ output_header = buildOutputHeader(input_headers.front(), params, true);
}
-
}
diff --git a/cpp-ch/local-engine/Operator/GraceMergingAggregatedStep.h
b/cpp-ch/local-engine/Operator/GraceMergingAggregatedStep.h
index d4d4fb41ed..94f6ba1af6 100644
--- a/cpp-ch/local-engine/Operator/GraceMergingAggregatedStep.h
+++ b/cpp-ch/local-engine/Operator/GraceMergingAggregatedStep.h
@@ -14,33 +14,25 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
+#pragma once
#include <unordered_map>
#include <Core/Block.h>
#include <Interpreters/Aggregator.h>
#include <Interpreters/Context.h>
-#include <Processors/Chunk.h>
-#include <Processors/IProcessor.h>
-#include <Processors/Port.h>
#include <Processors/QueryPlan/IQueryPlanStep.h>
#include <Processors/QueryPlan/ITransformingStep.h>
-#include <Processors/Transforms/AggregatingTransform.h>
-#include <QueryPipeline/SizeLimits.h>
-#include <Poco/Logger.h>
-#include <Common/AggregateUtil.h>
-#include <Common/logger_useful.h>
-#include "GraceAggregatingTransform.h"
namespace local_engine
{
/// It's used to merged aggregated data from intermediate aggregate stages,
it's the final stage of
/// aggregating.
-/// It support spilling data into disk when the memory usage is overflow.
+/// It supports spilling data into disk when the memory usage is overflow.
class GraceMergingAggregatedStep : public DB::ITransformingStep
{
public:
explicit GraceMergingAggregatedStep(
DB::ContextPtr context_,
- const DB::DataStream & input_stream_,
+ const DB::Block & input_header,
DB::Aggregator::Params params_,
bool no_pre_aggregated_);
~GraceMergingAggregatedStep() override = default;
@@ -55,7 +47,8 @@ private:
DB::ContextPtr context;
DB::Aggregator::Params params;
bool no_pre_aggregated;
- void updateOutputStream() override;
+protected:
+ void updateOutputHeader() override;
};
diff --git a/cpp-ch/local-engine/Operator/ReplicateRowsStep.cpp
b/cpp-ch/local-engine/Operator/ReplicateRowsStep.cpp
index ecb027c18f..705e413675 100644
--- a/cpp-ch/local-engine/Operator/ReplicateRowsStep.cpp
+++ b/cpp-ch/local-engine/Operator/ReplicateRowsStep.cpp
@@ -42,8 +42,8 @@ static DB::ITransformingStep::Traits getTraits()
}};
}
-ReplicateRowsStep::ReplicateRowsStep(const DB::DataStream & input_stream)
- : ITransformingStep(input_stream, transformHeader(input_stream.header),
getTraits())
+ReplicateRowsStep::ReplicateRowsStep(const DB::Block& input_header)
+ : ITransformingStep(input_header, transformHeader(input_header),
getTraits())
{
}
@@ -62,9 +62,9 @@ void
ReplicateRowsStep::transformPipeline(DB::QueryPipelineBuilder & pipeline, c
pipeline.addSimpleTransform([&](const DB::Block & header) { return
std::make_shared<ReplicateRowsTransform>(header); });
}
-void ReplicateRowsStep::updateOutputStream()
+void ReplicateRowsStep::updateOutputHeader()
{
- output_stream = createOutputStream(input_streams.front(),
transformHeader(input_streams.front().header), getDataStreamTraits());
+ output_header = transformHeader(input_headers.front());
}
ReplicateRowsTransform::ReplicateRowsTransform(const DB::Block & input_header_)
diff --git a/cpp-ch/local-engine/Operator/ReplicateRowsStep.h
b/cpp-ch/local-engine/Operator/ReplicateRowsStep.h
index f588bf0ceb..f4889e0907 100644
--- a/cpp-ch/local-engine/Operator/ReplicateRowsStep.h
+++ b/cpp-ch/local-engine/Operator/ReplicateRowsStep.h
@@ -25,24 +25,23 @@ namespace local_engine
class ReplicateRowsStep : public DB::ITransformingStep
{
public:
- ReplicateRowsStep(const DB::DataStream& input_stream);
+ ReplicateRowsStep(const DB::Block & input_header);
- static DB::Block transformHeader(const DB::Block& input);
+ static DB::Block transformHeader(const DB::Block & input);
String getName() const override { return "ReplicateRowsStep"; }
- void transformPipeline(DB::QueryPipelineBuilder& pipeline,
- const DB::BuildQueryPipelineSettings& settings)
override;
+ void transformPipeline(DB::QueryPipelineBuilder & pipeline, const
DB::BuildQueryPipelineSettings & settings) override;
+
private:
- void updateOutputStream() override;
+ void updateOutputHeader() override;
};
class ReplicateRowsTransform : public DB::ISimpleTransform
{
public:
- ReplicateRowsTransform(const DB::Block& input_header_);
+ ReplicateRowsTransform(const DB::Block & input_header_);
String getName() const override { return "ReplicateRowsTransform"; }
- void transform(DB::Chunk&) override;
-
+ void transform(DB::Chunk &) override;
};
}
diff --git a/cpp-ch/local-engine/Operator/StreamingAggregatingStep.cpp
b/cpp-ch/local-engine/Operator/StreamingAggregatingStep.cpp
index 2235f4cbe4..d19c9717d4 100644
--- a/cpp-ch/local-engine/Operator/StreamingAggregatingStep.cpp
+++ b/cpp-ch/local-engine/Operator/StreamingAggregatingStep.cpp
@@ -19,10 +19,10 @@
#include <Processors/Transforms/AggregatingTransform.h>
#include <QueryPipeline/QueryPipelineBuilder.h>
#include <Common/CHUtil.h>
-#include <Common/formatReadable.h>
#include <Common/GlutenConfig.h>
#include <Common/QueryContext.h>
#include <Common/Stopwatch.h>
+#include <Common/formatReadable.h>
namespace DB
{
@@ -34,7 +34,8 @@ namespace ErrorCodes
namespace local_engine
{
-StreamingAggregatingTransform::StreamingAggregatingTransform(DB::ContextPtr
context_, const DB::Block &header_, DB::AggregatingTransformParamsPtr params_)
+StreamingAggregatingTransform::StreamingAggregatingTransform(
+ DB::ContextPtr context_, const DB::Block & header_,
DB::AggregatingTransformParamsPtr params_)
: DB::IProcessor({header_}, {params_->getHeader()})
, context(context_)
, header(header_)
@@ -149,7 +150,7 @@ bool StreamingAggregatingTransform::needEvict()
/// If the grouping keys is high cardinality, we should evict data
variants early, and avoid to use a big
/// hash table.
- if (static_cast<double>(total_output_rows)/total_input_rows >
high_cardinality_threshold)
+ if (static_cast<double>(total_output_rows) / total_input_rows >
high_cardinality_threshold)
return true;
auto current_mem_used = currentThreadGroupMemoryUsage();
@@ -263,16 +264,14 @@ void StreamingAggregatingTransform::work()
static DB::ITransformingStep::Traits getTraits()
{
- return DB::ITransformingStep::Traits
- {
+ return DB::ITransformingStep::Traits{
{
.preserves_number_of_streams = false,
.preserves_sorting = false,
},
{
.preserves_number_of_rows = false,
- }
- };
+ }};
}
static DB::Block buildOutputHeader(const DB::Block & input_header_, const
DB::Aggregator::Params params_)
@@ -280,8 +279,8 @@ static DB::Block buildOutputHeader(const DB::Block &
input_header_, const DB::Ag
return params_.getHeader(input_header_, false);
}
StreamingAggregatingStep::StreamingAggregatingStep(
- DB::ContextPtr context_, const DB::DataStream & input_stream_,
DB::Aggregator::Params params_)
- : DB::ITransformingStep(input_stream_,
buildOutputHeader(input_stream_.header, params_), getTraits())
+ const DB::ContextPtr & context_, const DB::Block & input_header,
DB::Aggregator::Params params_)
+ : DB::ITransformingStep(input_header, buildOutputHeader(input_header,
params_), getTraits())
, context(context_)
, params(std::move(params_))
{
@@ -291,7 +290,8 @@ void
StreamingAggregatingStep::transformPipeline(DB::QueryPipelineBuilder & pipe
{
if (params.max_bytes_before_external_group_by)
{
- throw DB::Exception(DB::ErrorCodes::LOGICAL_ERROR,
"max_bytes_before_external_group_by is not supported in
StreamingAggregatingStep");
+ throw DB::Exception(
+ DB::ErrorCodes::LOGICAL_ERROR, "max_bytes_before_external_group_by
is not supported in StreamingAggregatingStep");
}
pipeline.dropTotalsAndExtremes();
auto transform_params =
std::make_shared<DB::AggregatingTransformParams>(pipeline.getHeader(), params,
false);
@@ -312,7 +312,7 @@ void
StreamingAggregatingStep::transformPipeline(DB::QueryPipelineBuilder & pipe
void
StreamingAggregatingStep::describeActions(DB::IQueryPlanStep::FormatSettings &
settings) const
{
- return params.explain(settings.out, settings.offset);
+ params.explain(settings.out, settings.offset);
}
void StreamingAggregatingStep::describeActions(DB::JSONBuilder::JSONMap & map)
const
@@ -320,9 +320,9 @@ void
StreamingAggregatingStep::describeActions(DB::JSONBuilder::JSONMap & map) c
params.explain(map);
}
-void StreamingAggregatingStep::updateOutputStream()
+void StreamingAggregatingStep::updateOutputHeader()
{
- output_stream = createOutputStream(input_streams.front(),
buildOutputHeader(input_streams.front().header, params), getDataStreamTraits());
+ output_header = buildOutputHeader(input_headers.front(), params);
}
}
diff --git a/cpp-ch/local-engine/Operator/StreamingAggregatingStep.h
b/cpp-ch/local-engine/Operator/StreamingAggregatingStep.h
index 0472d512c6..f3acb6be9b 100644
--- a/cpp-ch/local-engine/Operator/StreamingAggregatingStep.h
+++ b/cpp-ch/local-engine/Operator/StreamingAggregatingStep.h
@@ -16,19 +16,15 @@
*/
#pragma once
#include <cstddef>
-#include <unordered_map>
#include <Core/Block.h>
#include <Interpreters/Aggregator.h>
#include <Interpreters/Context.h>
#include <Processors/Chunk.h>
#include <Processors/IProcessor.h>
-#include <Processors/Port.h>
#include <Processors/QueryPlan/IQueryPlanStep.h>
#include <Processors/QueryPlan/ITransformingStep.h>
#include <Processors/Transforms/AggregatingTransform.h>
-#include <QueryPipeline/SizeLimits.h>
#include <Poco/Logger.h>
-#include <Common/logger_useful.h>
#include <Common/AggregateUtil.h>
namespace local_engine
@@ -40,11 +36,12 @@ class StreamingAggregatingTransform : public DB::IProcessor
{
public:
using Status = DB::IProcessor::Status;
- explicit StreamingAggregatingTransform(DB::ContextPtr context_, const
DB::Block &header_, DB::AggregatingTransformParamsPtr params_);
+ explicit StreamingAggregatingTransform(DB::ContextPtr context_, const
DB::Block & header_, DB::AggregatingTransformParamsPtr params_);
~StreamingAggregatingTransform() override;
String getName() const override { return "StreamingAggregatingTransform"; }
Status prepare() override;
void work() override;
+
private:
DB::ContextPtr context;
DB::Block header;
@@ -71,7 +68,7 @@ private:
DB::Chunk input_chunk;
DB::Chunk output_chunk;
bool input_finished = false;
-
+
std::unique_ptr<AggregateDataBlockConverter> block_converter = nullptr;
Poco::Logger * logger =
&Poco::Logger::get("StreamingAggregatingTransform");
@@ -92,10 +89,7 @@ private:
class StreamingAggregatingStep : public DB::ITransformingStep
{
public:
- explicit StreamingAggregatingStep(
- DB::ContextPtr context_,
- const DB::DataStream & input_stream_,
- DB::Aggregator::Params params_);
+ explicit StreamingAggregatingStep(const DB::ContextPtr & context_, const
DB::Block & input_header, DB::Aggregator::Params params_);
~StreamingAggregatingStep() override = default;
String getName() const override { return "StreamingAggregating"; }
@@ -108,6 +102,8 @@ public:
private:
DB::ContextPtr context;
DB::Aggregator::Params params;
- void updateOutputStream() override;
+
+protected:
+ void updateOutputHeader() override;
};
}
diff --git a/cpp-ch/local-engine/Operator/WindowGroupLimitStep.cpp
b/cpp-ch/local-engine/Operator/WindowGroupLimitStep.cpp
index d2bd583a1c..f25e3f22ac 100644
--- a/cpp-ch/local-engine/Operator/WindowGroupLimitStep.cpp
+++ b/cpp-ch/local-engine/Operator/WindowGroupLimitStep.cpp
@@ -302,12 +302,12 @@ static DB::ITransformingStep::Traits getTraits()
}
WindowGroupLimitStep::WindowGroupLimitStep(
- const DB::DataStream & input_stream_,
+ const DB::Block & input_header_,
const String & function_name_,
- const std::vector<size_t> partition_columns_,
- const std::vector<size_t> sort_columns_,
+ const std::vector<size_t> & partition_columns_,
+ const std::vector<size_t> & sort_columns_,
size_t limit_)
- : DB::ITransformingStep(input_stream_, input_stream_.header, getTraits())
+ : DB::ITransformingStep(input_header_, input_header_, getTraits())
, function_name(function_name_)
, partition_columns(partition_columns_)
, sort_columns(sort_columns_)
@@ -321,9 +321,9 @@ void
WindowGroupLimitStep::describePipeline(DB::IQueryPlanStep::FormatSettings &
DB::IQueryPlanStep::describePipeline(processors, settings);
}
-void WindowGroupLimitStep::updateOutputStream()
+void WindowGroupLimitStep::updateOutputHeader()
{
- output_stream = createOutputStream(input_streams.front(),
input_streams.front().header, getDataStreamTraits());
+ output_header = input_headers.front();
}
diff --git a/cpp-ch/local-engine/Operator/WindowGroupLimitStep.h
b/cpp-ch/local-engine/Operator/WindowGroupLimitStep.h
index bbbbf42abc..55e3eaeb72 100644
--- a/cpp-ch/local-engine/Operator/WindowGroupLimitStep.h
+++ b/cpp-ch/local-engine/Operator/WindowGroupLimitStep.h
@@ -27,10 +27,10 @@ class WindowGroupLimitStep : public DB::ITransformingStep
{
public:
explicit WindowGroupLimitStep(
- const DB::DataStream & input_stream_,
+ const DB::Block & input_header_,
const String & function_name_,
- const std::vector<size_t> partition_columns_,
- const std::vector<size_t> sort_columns_,
+ const std::vector<size_t> & partition_columns_,
+ const std::vector<size_t> & sort_columns_,
size_t limit_);
~WindowGroupLimitStep() override = default;
@@ -38,7 +38,7 @@ public:
void transformPipeline(DB::QueryPipelineBuilder & pipeline, const
DB::BuildQueryPipelineSettings & settings) override;
void describePipeline(DB::IQueryPlanStep::FormatSettings & settings) const
override;
- void updateOutputStream() override;
+ void updateOutputHeader() override;
private:
// window function name, one of row_number, rank and dense_rank
diff --git a/cpp-ch/local-engine/Parser/InputFileNameParser.cpp
b/cpp-ch/local-engine/Parser/InputFileNameParser.cpp
index 8fcf746112..e6fb2fa6b3 100644
--- a/cpp-ch/local-engine/Parser/InputFileNameParser.cpp
+++ b/cpp-ch/local-engine/Parser/InputFileNameParser.cpp
@@ -40,13 +40,13 @@ static DB::ITransformingStep::Traits getTraits()
}};
}
-static DB::Block getOutputHeader(
- const DB::DataStream & input_stream,
+static DB::Block createOutputHeader(
+ const DB::Block & header,
const std::optional<String> & file_name,
const std::optional<Int64> & block_start,
const std::optional<Int64> & block_length)
{
- DB::Block output_header = input_stream.header;
+ DB::Block output_header{header};
if (file_name.has_value())
output_header.insert(DB::ColumnWithTypeAndName{std::make_shared<DB::DataTypeString>(),
InputFileNameParser::INPUT_FILE_NAME});
if (block_start.has_value())
@@ -89,11 +89,11 @@ class InputFileExprProjectStep : public
DB::ITransformingStep
{
public:
InputFileExprProjectStep(
- const DB::DataStream & input_stream,
+ const DB::Block & input_header,
const std::optional<String> & file_name,
const std::optional<Int64> & block_start,
const std::optional<Int64> & block_length)
- : ITransformingStep(input_stream, getOutputHeader(input_stream,
file_name, block_start, block_length), getTraits(), true)
+ : ITransformingStep(input_header, createOutputHeader(input_header,
file_name, block_start, block_length), getTraits(), true)
, file_name(file_name)
, block_start(block_start)
, block_length(block_length)
@@ -105,15 +105,14 @@ public:
void transformPipeline(DB::QueryPipelineBuilder & pipeline, const
DB::BuildQueryPipelineSettings & /*settings*/) override
{
pipeline.addSimpleTransform(
- [&](const DB::Block & header) {
- return std::make_shared<InputFileExprProjectTransform>(header,
output_stream->header, file_name, block_start, block_length);
- });
+ [&](const DB::Block & header)
+ { return std::make_shared<InputFileExprProjectTransform>(header,
*output_header, file_name, block_start, block_length); });
}
protected:
- void updateOutputStream() override
+ void updateOutputHeader() override
{
- output_stream = createOutputStream(input_streams.front(),
output_stream->header, getDataStreamTraits());
+ // do nothing
}
private:
@@ -198,14 +197,14 @@ std::optional<DB::IQueryPlanStep *>
InputFileNameParser::addInputFileProjectStep
{
if (!file_name.has_value() && !block_start.has_value() &&
!block_length.has_value())
return std::nullopt;
- auto step =
std::make_unique<InputFileExprProjectStep>(plan.getCurrentDataStream(),
file_name, block_start, block_length);
+ auto step =
std::make_unique<InputFileExprProjectStep>(plan.getCurrentHeader(), file_name,
block_start, block_length);
step->setStepDescription("Input file expression project");
std::optional<DB::IQueryPlanStep *> result = step.get();
plan.addStep(std::move(step));
return result;
}
-void InputFileNameParser::addInputFileColumnsToChunk(const DB::Block & header,
DB::Chunk & chunk)
+void InputFileNameParser::addInputFileColumnsToChunk(const DB::Block & header,
DB::Chunk & chunk) const
{
addInputFileColumnsToChunk(header, chunk, file_name, block_start,
block_length);
}
diff --git a/cpp-ch/local-engine/Parser/InputFileNameParser.h
b/cpp-ch/local-engine/Parser/InputFileNameParser.h
index 09b3e72617..ab52ba1072 100644
--- a/cpp-ch/local-engine/Parser/InputFileNameParser.h
+++ b/cpp-ch/local-engine/Parser/InputFileNameParser.h
@@ -52,7 +52,7 @@ public:
void setBlockLength(const Int64 block_length) { this->block_length =
block_length; }
[[nodiscard]] std::optional<DB::IQueryPlanStep *>
addInputFileProjectStep(DB::QueryPlan & plan);
- void addInputFileColumnsToChunk(const DB::Block & header, DB::Chunk &
chunk);
+ void addInputFileColumnsToChunk(const DB::Block & header, DB::Chunk &
chunk) const;
private:
std::optional<String> file_name;
diff --git a/cpp-ch/local-engine/Parser/LocalExecutor.cpp
b/cpp-ch/local-engine/Parser/LocalExecutor.cpp
index f781556c91..ebbf5064c9 100644
--- a/cpp-ch/local-engine/Parser/LocalExecutor.cpp
+++ b/cpp-ch/local-engine/Parser/LocalExecutor.cpp
@@ -137,7 +137,7 @@ Block LocalExecutor::getHeader()
LocalExecutor::LocalExecutor(QueryPlanPtr query_plan, QueryPipelineBuilderPtr
pipeline_builder, bool dump_pipeline_)
: query_pipeline_builder(std::move(pipeline_builder))
- , header(query_plan->getCurrentDataStream().header.cloneEmpty())
+ , header(query_plan->getCurrentHeader().cloneEmpty())
, dump_pipeline(dump_pipeline_)
, ch_column_to_spark_row(std::make_unique<CHColumnToSparkRow>())
, current_query_plan(std::move(query_plan))
diff --git a/cpp-ch/local-engine/Parser/RelParsers/AggregateRelParser.cpp
b/cpp-ch/local-engine/Parser/RelParsers/AggregateRelParser.cpp
index 6875b10d13..adde9ac182 100644
--- a/cpp-ch/local-engine/Parser/RelParsers/AggregateRelParser.cpp
+++ b/cpp-ch/local-engine/Parser/RelParsers/AggregateRelParser.cpp
@@ -76,32 +76,32 @@ AggregateRelParser::parse(DB::QueryPlanPtr query_plan,
const substrait::Rel & re
setup(std::move(query_plan), rel);
addPreProjection();
- LOG_TRACE(logger, "header after pre-projection is: {}",
plan->getCurrentDataStream().header.dumpStructure());
+ LOG_TRACE(logger, "header after pre-projection is: {}",
plan->getCurrentHeader().dumpStructure());
if (has_final_stage)
{
addMergingAggregatedStep();
- LOG_TRACE(logger, "header after merging is: {}",
plan->getCurrentDataStream().header.dumpStructure());
+ LOG_TRACE(logger, "header after merging is: {}",
plan->getCurrentHeader().dumpStructure());
addPostProjection();
- LOG_TRACE(logger, "header after post-projection is: {}",
plan->getCurrentDataStream().header.dumpStructure());
+ LOG_TRACE(logger, "header after post-projection is: {}",
plan->getCurrentHeader().dumpStructure());
}
else if (has_complete_stage)
{
addCompleteModeAggregatedStep();
- LOG_TRACE(logger, "header after complete aggregate is: {}",
plan->getCurrentDataStream().header.dumpStructure());
+ LOG_TRACE(logger, "header after complete aggregate is: {}",
plan->getCurrentHeader().dumpStructure());
addPostProjection();
- LOG_TRACE(logger, "header after post-projection is: {}",
plan->getCurrentDataStream().header.dumpStructure());
+ LOG_TRACE(logger, "header after post-projection is: {}",
plan->getCurrentHeader().dumpStructure());
}
else
{
addAggregatingStep();
- LOG_TRACE(logger, "header after aggregating is: {}",
plan->getCurrentDataStream().header.dumpStructure());
+ LOG_TRACE(logger, "header after aggregating is: {}",
plan->getCurrentHeader().dumpStructure());
}
/// Add a check here to help find bugs, Don't remove it.
/// Thre order of result of result columns must be ordered grouping keys
++ ordered aggregate expression results
- auto aggregation_output_header = plan->getCurrentDataStream().header;
+ auto aggregation_output_header = plan->getCurrentHeader();
for (size_t i = 0; i < grouping_keys.size(); ++i)
{
auto pos =
aggregation_output_header.getPositionByName(grouping_keys[i]);
@@ -114,7 +114,7 @@ AggregateRelParser::parse(DB::QueryPlanPtr query_plan,
const substrait::Rel & re
&& (has_final_stage || has_complete_stage ||
rel.aggregate().measures().empty()))
{
LOG_TRACE(&Poco::Logger::get("AggregateRelParser"), "default aggregate
result step");
- auto default_agg_result =
std::make_unique<DefaultHashAggregateResultStep>(plan->getCurrentDataStream());
+ auto default_agg_result =
std::make_unique<DefaultHashAggregateResultStep>(plan->getCurrentHeader());
default_agg_result->setStepDescription("Default aggregate result");
steps.push_back(default_agg_result.get());
plan->addStep(std::move(default_agg_result));
@@ -176,7 +176,7 @@ void AggregateRelParser::setup(DB::QueryPlanPtr query_plan,
const substrait::Rel
DB::ErrorCodes::LOGICAL_ERROR, "AggregateRelParser: multiple
aggregation phases with complete mode are not supported");
}
- auto input_header = plan->getCurrentDataStream().header;
+ auto input_header = plan->getCurrentHeader();
for (const auto & measure : aggregate_rel->measures())
{
AggregateInfo agg_info;
@@ -219,7 +219,7 @@ void AggregateRelParser::setup(DB::QueryPlanPtr query_plan,
const substrait::Rel
/// The projections are built by the function parsers.
void AggregateRelParser::addPreProjection()
{
- auto input_header = plan->getCurrentDataStream().header;
+ auto input_header = plan->getCurrentHeader();
ActionsDAG projection_action{input_header.getColumnsWithTypeAndName()};
std::string dag_footprint = projection_action.dumpDAG();
for (auto & agg_info : aggregates)
@@ -238,7 +238,7 @@ void AggregateRelParser::addPreProjection()
{
/// Avoid unnecessary evaluation
projection_action.removeUnusedActions();
- auto projection_step =
std::make_unique<DB::ExpressionStep>(plan->getCurrentDataStream(),
std::move(projection_action));
+ auto projection_step =
std::make_unique<DB::ExpressionStep>(plan->getCurrentHeader(),
std::move(projection_action));
projection_step->setStepDescription("Projection before aggregate");
steps.emplace_back(projection_step.get());
plan->addStep(std::move(projection_step));
@@ -247,7 +247,7 @@ void AggregateRelParser::addPreProjection()
void AggregateRelParser::buildAggregateDescriptions(AggregateDescriptions &
descriptions)
{
- const auto & current_plan_header = plan->getCurrentDataStream().header;
+ const auto & current_plan_header = plan->getCurrentHeader();
auto build_result_column_name
= [this, current_plan_header](
const String & function_name, const Array & params, const
Strings & arg_names, substrait::AggregationPhase phase)
@@ -356,7 +356,7 @@ void AggregateRelParser::addMergingAggregatedStep()
if (config.enable_streaming_aggregating)
{
params.group_by_two_level_threshold =
settings[Setting::group_by_two_level_threshold];
- auto merging_step =
std::make_unique<GraceMergingAggregatedStep>(getContext(),
plan->getCurrentDataStream(), params, false);
+ auto merging_step =
std::make_unique<GraceMergingAggregatedStep>(getContext(),
plan->getCurrentHeader(), params, false);
steps.emplace_back(merging_step.get());
plan->addStep(std::move(merging_step));
}
@@ -365,7 +365,7 @@ void AggregateRelParser::addMergingAggregatedStep()
/// We don't use the grouping set feature in CH, so
grouping_sets_params_list should always be empty.
DB::GroupingSetsParamsList grouping_sets_params_list;
auto merging_step = std::make_unique<DB::MergingAggregatedStep>(
- plan->getCurrentDataStream(),
+ plan->getCurrentHeader(),
params,
grouping_sets_params_list,
true,
@@ -410,7 +410,7 @@ void AggregateRelParser::addCompleteModeAggregatedStep()
settings[Setting::optimize_group_by_constant_keys],
settings[Setting::min_hit_rate_to_use_consecutive_keys_optimization],
/*StatsCollectingParams*/ {});
- auto merging_step =
std::make_unique<GraceMergingAggregatedStep>(getContext(),
plan->getCurrentDataStream(), params, true);
+ auto merging_step =
std::make_unique<GraceMergingAggregatedStep>(getContext(),
plan->getCurrentHeader(), params, true);
steps.emplace_back(merging_step.get());
plan->addStep(std::move(merging_step));
}
@@ -439,7 +439,7 @@ void AggregateRelParser::addCompleteModeAggregatedStep()
/*StatsCollectingParams*/ {});
auto aggregating_step = std::make_unique<AggregatingStep>(
- plan->getCurrentDataStream(),
+ plan->getCurrentHeader(),
params,
GroupingSetsParamsList(),
true,
@@ -507,7 +507,7 @@ void AggregateRelParser::addAggregatingStep()
/*StatsCollectingParams*/ {});
if (!is_distinct_aggreate)
{
- auto aggregating_step =
std::make_unique<StreamingAggregatingStep>(getContext(),
plan->getCurrentDataStream(), params);
+ auto aggregating_step =
std::make_unique<StreamingAggregatingStep>(getContext(),
plan->getCurrentHeader(), params);
steps.emplace_back(aggregating_step.get());
plan->addStep(std::move(aggregating_step));
}
@@ -526,7 +526,7 @@ void AggregateRelParser::addAggregatingStep()
/// will make the result for count(distinct(n_name)) wrong. step3
must finish all inputs before it puts any block into step4.
/// So we introduce GraceAggregatingStep here, it can handle mass
data with high cardinality.
auto aggregating_step
- = std::make_unique<GraceAggregatingStep>(getContext(),
plan->getCurrentDataStream(), params, has_first_stage);
+ = std::make_unique<GraceAggregatingStep>(getContext(),
plan->getCurrentHeader(), params, has_first_stage);
steps.emplace_back(aggregating_step.get());
plan->addStep(std::move(aggregating_step));
}
@@ -556,7 +556,7 @@ void AggregateRelParser::addAggregatingStep()
/*StatsCollectingParams*/ {});
auto aggregating_step = std::make_unique<AggregatingStep>(
- plan->getCurrentDataStream(),
+ plan->getCurrentHeader(),
params,
GroupingSetsParamsList(),
false,
@@ -579,7 +579,7 @@ void AggregateRelParser::addAggregatingStep()
// Only be called in final stage.
void AggregateRelParser::addPostProjection()
{
- auto input_header = plan->getCurrentDataStream().header;
+ auto input_header = plan->getCurrentHeader();
ActionsDAG project_actions_dag{input_header.getColumnsWithTypeAndName()};
auto dag_footprint = project_actions_dag.dumpDAG();
@@ -612,7 +612,7 @@ void AggregateRelParser::addPostProjection()
}
if (project_actions_dag.dumpDAG() != dag_footprint)
{
- QueryPlanStepPtr convert_step =
std::make_unique<ExpressionStep>(plan->getCurrentDataStream(),
std::move(project_actions_dag));
+ QueryPlanStepPtr convert_step =
std::make_unique<ExpressionStep>(plan->getCurrentHeader(),
std::move(project_actions_dag));
convert_step->setStepDescription("Post-projection for aggregate");
steps.emplace_back(convert_step.get());
plan->addStep(std::move(convert_step));
diff --git a/cpp-ch/local-engine/Parser/RelParsers/CrossRelParser.cpp
b/cpp-ch/local-engine/Parser/RelParsers/CrossRelParser.cpp
index 0620a87fb0..4fef282fe4 100644
--- a/cpp-ch/local-engine/Parser/RelParsers/CrossRelParser.cpp
+++ b/cpp-ch/local-engine/Parser/RelParsers/CrossRelParser.cpp
@@ -98,12 +98,12 @@ void CrossRelParser::renamePlanColumns(DB::QueryPlan &
left, DB::QueryPlan & rig
{
/// To support mixed join conditions, we must make sure that the column
names in the right be the same as
/// storage_join's right sample block.
- auto right_ori_header =
right.getCurrentDataStream().header.getColumnsWithTypeAndName();
+ auto right_ori_header =
right.getCurrentHeader().getColumnsWithTypeAndName();
if (right_ori_header.size() > 0 && right_ori_header[0].name !=
BlockUtil::VIRTUAL_ROW_COUNT_COLUMN)
{
ActionsDAG right_project = ActionsDAG::makeConvertingActions(
right_ori_header,
storage_join.getRightSampleBlock().getColumnsWithTypeAndName(),
ActionsDAG::MatchColumnsMode::Position);
- QueryPlanStepPtr project_step =
std::make_unique<ExpressionStep>(right.getCurrentDataStream(),
std::move(right_project));
+ QueryPlanStepPtr project_step =
std::make_unique<ExpressionStep>(right.getCurrentHeader(),
std::move(right_project));
project_step->setStepDescription("Rename Broadcast Table Name");
steps.emplace_back(project_step.get());
right.addStep(std::move(project_step));
@@ -113,17 +113,17 @@ void CrossRelParser::renamePlanColumns(DB::QueryPlan &
left, DB::QueryPlan & rig
/// avoid the columns name in the right table be changed in
`addConvertStep`.
/// This could happen in tpc-ds q44.
DB::ColumnsWithTypeAndName new_left_cols;
- const auto & right_header = right.getCurrentDataStream().header;
+ const auto & right_header = right.getCurrentHeader();
auto left_prefix = getUniqueName("left");
- for (const auto & col : left.getCurrentDataStream().header)
+ for (const auto & col : left.getCurrentHeader())
if (right_header.has(col.name))
new_left_cols.emplace_back(col.column, col.type, left_prefix +
col.name);
else
new_left_cols.emplace_back(col.column, col.type, col.name);
- auto left_header =
left.getCurrentDataStream().header.getColumnsWithTypeAndName();
+ auto left_header = left.getCurrentHeader().getColumnsWithTypeAndName();
ActionsDAG left_project = ActionsDAG::makeConvertingActions(left_header,
new_left_cols, ActionsDAG::MatchColumnsMode::Position);
- QueryPlanStepPtr project_step =
std::make_unique<ExpressionStep>(left.getCurrentDataStream(),
std::move(left_project));
+ QueryPlanStepPtr project_step =
std::make_unique<ExpressionStep>(left.getCurrentHeader(),
std::move(left_project));
project_step->setStepDescription("Rename Left Table Name for broadcast
join");
steps.emplace_back(project_step.get());
left.addStep(std::move(project_step));
@@ -139,28 +139,28 @@ DB::QueryPlanPtr CrossRelParser::parseJoin(const
substrait::CrossRel & join, DB:
if (storage_join)
renamePlanColumns(*left, *right, *storage_join);
auto table_join = createCrossTableJoin(join.type());
- DB::Block right_header_before_convert_step =
right->getCurrentDataStream().header;
+ DB::Block right_header_before_convert_step = right->getCurrentHeader();
addConvertStep(*table_join, *left, *right);
// Add a check to find error easily.
- if (!blocksHaveEqualStructure(right_header_before_convert_step,
right->getCurrentDataStream().header))
+ if (!blocksHaveEqualStructure(right_header_before_convert_step,
right->getCurrentHeader()))
{
throw DB::Exception(
DB::ErrorCodes::LOGICAL_ERROR,
"For broadcast join, we must not change the columns name in the
right table.\nleft header:{},\nright header: {} -> {}",
- left->getCurrentDataStream().header.dumpNames(),
+ left->getCurrentHeader().dumpNames(),
right_header_before_convert_step.dumpNames(),
- right->getCurrentDataStream().header.dumpNames());
+ right->getCurrentHeader().dumpNames());
}
Names after_join_names;
- auto left_names = left->getCurrentDataStream().header.getNames();
+ auto left_names = left->getCurrentHeader().getNames();
after_join_names.insert(after_join_names.end(), left_names.begin(),
left_names.end());
auto right_name = table_join->columnsFromJoinedTable().getNames();
after_join_names.insert(after_join_names.end(), right_name.begin(),
right_name.end());
- auto left_header = left->getCurrentDataStream().header;
- auto right_header = right->getCurrentDataStream().header;
+ auto left_header = left->getCurrentHeader();
+ auto right_header = right->getCurrentHeader();
QueryPlanPtr query_plan;
if (storage_join)
@@ -171,7 +171,7 @@ DB::QueryPlanPtr CrossRelParser::parseJoin(const
substrait::CrossRel & join, DB:
auto broadcast_hash_join = storage_join->getJoinLocked(table_join,
context);
// table_join->resetKeys();
- QueryPlanStepPtr join_step =
std::make_unique<FilledJoinStep>(left->getCurrentDataStream(),
broadcast_hash_join, 8192);
+ QueryPlanStepPtr join_step =
std::make_unique<FilledJoinStep>(left->getCurrentHeader(), broadcast_hash_join,
8192);
join_step->setStepDescription("STORAGE_JOIN");
steps.emplace_back(join_step.get());
@@ -193,9 +193,9 @@ DB::QueryPlanPtr CrossRelParser::parseJoin(const
substrait::CrossRel & join, DB:
}
else
{
- JoinPtr hash_join = std::make_shared<HashJoin>(table_join,
right->getCurrentDataStream().header.cloneEmpty());
+ JoinPtr hash_join = std::make_shared<HashJoin>(table_join,
right->getCurrentHeader().cloneEmpty());
QueryPlanStepPtr join_step
- = std::make_unique<DB::JoinStep>(left->getCurrentDataStream(),
right->getCurrentDataStream(), hash_join, 8192, 1, false);
+ = std::make_unique<DB::JoinStep>(left->getCurrentHeader(),
right->getCurrentHeader(), hash_join, 8192, 1, false);
join_step->setStepDescription("CROSS_JOIN");
steps.emplace_back(join_step.get());
std::vector<QueryPlanPtr> plans;
@@ -218,7 +218,7 @@ void CrossRelParser::addPostFilter(DB::QueryPlan &
query_plan, const substrait::
auto expression = join_rel.expression();
std::string filter_name;
- ActionsDAG
actions_dag(query_plan.getCurrentDataStream().header.getColumnsWithTypeAndName());
+ ActionsDAG
actions_dag(query_plan.getCurrentHeader().getColumnsWithTypeAndName());
if (!expression.has_scalar_function())
{
// It may be singular_or_list
@@ -230,7 +230,7 @@ void CrossRelParser::addPostFilter(DB::QueryPlan &
query_plan, const substrait::
const auto * func_node =
expression_parser->parseFunction(expression.scalar_function(), actions_dag,
true);
filter_name = func_node->result_name;
}
- auto filter_step =
std::make_unique<FilterStep>(query_plan.getCurrentDataStream(),
std::move(actions_dag), filter_name, true);
+ auto filter_step =
std::make_unique<FilterStep>(query_plan.getCurrentHeader(),
std::move(actions_dag), filter_name, true);
filter_step->setStepDescription("Post Join Filter");
steps.emplace_back(filter_step.get());
query_plan.addStep(std::move(filter_step));
@@ -240,16 +240,16 @@ void CrossRelParser::addConvertStep(TableJoin &
table_join, DB::QueryPlan & left
{
/// If the columns name in right table is duplicated with left table, we
need to rename the right table's columns.
NameSet left_columns_set;
- for (const auto & col : left.getCurrentDataStream().header.getNames())
+ for (const auto & col : left.getCurrentHeader().getNames())
left_columns_set.emplace(col);
table_join.setColumnsFromJoinedTable(
- right.getCurrentDataStream().header.getNamesAndTypesList(),
left_columns_set, getUniqueName("right") + ".");
+ right.getCurrentHeader().getNamesAndTypesList(), left_columns_set,
getUniqueName("right") + ".");
// fix right table key duplicate
NamesWithAliases right_table_alias;
for (size_t idx = 0; idx < table_join.columnsFromJoinedTable().size();
idx++)
{
- auto origin_name =
right.getCurrentDataStream().header.getByPosition(idx).name;
+ auto origin_name = right.getCurrentHeader().getByPosition(idx).name;
auto dedup_name =
table_join.columnsFromJoinedTable().getNames().at(idx);
if (origin_name != dedup_name)
{
@@ -258,8 +258,8 @@ void CrossRelParser::addConvertStep(TableJoin & table_join,
DB::QueryPlan & left
}
if (!right_table_alias.empty())
{
- ActionsDAG
rename_dag(right.getCurrentDataStream().header.getNamesAndTypesList());
- auto original_right_columns = right.getCurrentDataStream().header;
+ ActionsDAG rename_dag(right.getCurrentHeader().getNamesAndTypesList());
+ auto original_right_columns = right.getCurrentHeader();
for (const auto & column_alias : right_table_alias)
{
if (original_right_columns.has(column_alias.first))
@@ -270,7 +270,7 @@ void CrossRelParser::addConvertStep(TableJoin & table_join,
DB::QueryPlan & left
}
}
- QueryPlanStepPtr project_step =
std::make_unique<ExpressionStep>(right.getCurrentDataStream(),
std::move(rename_dag));
+ QueryPlanStepPtr project_step =
std::make_unique<ExpressionStep>(right.getCurrentHeader(),
std::move(rename_dag));
project_step->setStepDescription("Right Table Rename");
steps.emplace_back(project_step.get());
right.addStep(std::move(project_step));
@@ -283,11 +283,11 @@ void CrossRelParser::addConvertStep(TableJoin &
table_join, DB::QueryPlan & left
std::optional<ActionsDAG> left_convert_actions;
std::optional<ActionsDAG> right_convert_actions;
std::tie(left_convert_actions, right_convert_actions) =
table_join.createConvertingActions(
- left.getCurrentDataStream().header.getColumnsWithTypeAndName(),
right.getCurrentDataStream().header.getColumnsWithTypeAndName());
+ left.getCurrentHeader().getColumnsWithTypeAndName(),
right.getCurrentHeader().getColumnsWithTypeAndName());
if (right_convert_actions)
{
- auto converting_step =
std::make_unique<ExpressionStep>(right.getCurrentDataStream(),
std::move(*right_convert_actions));
+ auto converting_step =
std::make_unique<ExpressionStep>(right.getCurrentHeader(),
std::move(*right_convert_actions));
converting_step->setStepDescription("Convert joined columns");
steps.emplace_back(converting_step.get());
right.addStep(std::move(converting_step));
@@ -295,7 +295,7 @@ void CrossRelParser::addConvertStep(TableJoin & table_join,
DB::QueryPlan & left
if (left_convert_actions)
{
- auto converting_step =
std::make_unique<ExpressionStep>(left.getCurrentDataStream(),
std::move(*left_convert_actions));
+ auto converting_step =
std::make_unique<ExpressionStep>(left.getCurrentHeader(),
std::move(*left_convert_actions));
converting_step->setStepDescription("Convert joined columns");
steps.emplace_back(converting_step.get());
left.addStep(std::move(converting_step));
diff --git a/cpp-ch/local-engine/Parser/RelParsers/ExpandRelParser.cpp
b/cpp-ch/local-engine/Parser/RelParsers/ExpandRelParser.cpp
index 16470ab64a..8a64c445e7 100644
--- a/cpp-ch/local-engine/Parser/RelParsers/ExpandRelParser.cpp
+++ b/cpp-ch/local-engine/Parser/RelParsers/ExpandRelParser.cpp
@@ -48,7 +48,7 @@ void updateType(DB::DataTypePtr & type, const DB::DataTypePtr
& new_type)
DB::QueryPlanPtr ExpandRelParser::parse(DB::QueryPlanPtr query_plan, const
substrait::Rel & rel, std::list<const substrait::Rel *> &)
{
const auto & expand_rel = rel.expand();
- const auto & header = query_plan->getCurrentDataStream().header;
+ const auto & header = query_plan->getCurrentHeader();
std::vector<std::vector<ExpandFieldKind>> expand_kinds;
std::vector<std::vector<DB::Field>> expand_fields;
@@ -123,7 +123,7 @@ DB::QueryPlanPtr ExpandRelParser::parse(DB::QueryPlanPtr
query_plan, const subst
}
ExpandField expand_field(names, types, expand_kinds, expand_fields);
- auto expand_step =
std::make_unique<ExpandStep>(query_plan->getCurrentDataStream(),
std::move(expand_field));
+ auto expand_step =
std::make_unique<ExpandStep>(query_plan->getCurrentHeader(),
std::move(expand_field));
expand_step->setStepDescription("Expand Step");
steps.emplace_back(expand_step.get());
query_plan->addStep(std::move(expand_step));
diff --git a/cpp-ch/local-engine/Parser/RelParsers/FetchRelParser.cpp
b/cpp-ch/local-engine/Parser/RelParsers/FetchRelParser.cpp
index 0497c8e371..1b4c5e037e 100644
--- a/cpp-ch/local-engine/Parser/RelParsers/FetchRelParser.cpp
+++ b/cpp-ch/local-engine/Parser/RelParsers/FetchRelParser.cpp
@@ -31,7 +31,7 @@ public:
DB::QueryPlanPtr parse(DB::QueryPlanPtr query_plan, const substrait::Rel &
rel, std::list<const substrait::Rel *> &)
{
const auto & limit = rel.fetch();
- auto limit_step =
std::make_unique<DB::LimitStep>(query_plan->getCurrentDataStream(),
limit.count(), limit.offset());
+ auto limit_step =
std::make_unique<DB::LimitStep>(query_plan->getCurrentHeader(), limit.count(),
limit.offset());
limit_step->setStepDescription("LIMIT");
steps.push_back(limit_step.get());
query_plan->addStep(std::move(limit_step));
diff --git a/cpp-ch/local-engine/Parser/RelParsers/FilterRelParser.cpp
b/cpp-ch/local-engine/Parser/RelParsers/FilterRelParser.cpp
index 21f83a0f2c..eff1033cb5 100644
--- a/cpp-ch/local-engine/Parser/RelParsers/FilterRelParser.cpp
+++ b/cpp-ch/local-engine/Parser/RelParsers/FilterRelParser.cpp
@@ -35,7 +35,7 @@ FilterRelParser::parse(DB::QueryPlanPtr query_plan, const
substrait::Rel & rel,
const auto & filter_rel = final_rel.filter();
std::string filter_name;
- auto input_header = query_plan->getCurrentDataStream().header;
+ auto input_header = query_plan->getCurrentHeader();
DB::ActionsDAG actions_dag{input_header.getColumnsWithTypeAndName()};
const auto condition_node = parseExpression(actions_dag,
filter_rel.condition());
if (filter_rel.condition().has_scalar_function())
@@ -45,7 +45,7 @@ FilterRelParser::parse(DB::QueryPlanPtr query_plan, const
substrait::Rel & rel,
filter_name = condition_node->result_name;
bool remove_filter_column = true;
- auto input_names = query_plan->getCurrentDataStream().header.getNames();
+ auto input_names = query_plan->getCurrentHeader().getNames();
DB::NameSet input_with_condition(input_names.begin(), input_names.end());
if (input_with_condition.contains(condition_node->result_name))
remove_filter_column = false;
@@ -57,13 +57,13 @@ FilterRelParser::parse(DB::QueryPlanPtr query_plan, const
substrait::Rel & rel,
auto non_nullable_columns = non_nullable_columns_resolver.resolve();
auto filter_step
- = std::make_unique<DB::FilterStep>(query_plan->getCurrentDataStream(),
std::move(actions_dag), filter_name, remove_filter_column);
+ = std::make_unique<DB::FilterStep>(query_plan->getCurrentHeader(),
std::move(actions_dag), filter_name, remove_filter_column);
filter_step->setStepDescription("WHERE");
steps.emplace_back(filter_step.get());
query_plan->addStep(std::move(filter_step));
// header maybe changed, need to rollback it
- if (!blocksHaveEqualStructure(input_header,
query_plan->getCurrentDataStream().header))
+ if (!blocksHaveEqualStructure(input_header,
query_plan->getCurrentHeader()))
{
steps.emplace_back(PlanUtil::adjustQueryPlanHeader(*query_plan,
input_header, "Rollback filter header"));
}
diff --git a/cpp-ch/local-engine/Parser/RelParsers/JoinRelParser.cpp
b/cpp-ch/local-engine/Parser/RelParsers/JoinRelParser.cpp
index 1157e8dbb0..c8b81b21a0 100644
--- a/cpp-ch/local-engine/Parser/RelParsers/JoinRelParser.cpp
+++ b/cpp-ch/local-engine/Parser/RelParsers/JoinRelParser.cpp
@@ -163,11 +163,11 @@ void JoinRelParser::renamePlanColumns(DB::QueryPlan &
left, DB::QueryPlan & righ
/// To support mixed join conditions, we must make sure that the column
names in the right be the same as
/// storage_join's right sample block.
ActionsDAG right_project = ActionsDAG::makeConvertingActions(
- right.getCurrentDataStream().header.getColumnsWithTypeAndName(),
+ right.getCurrentHeader().getColumnsWithTypeAndName(),
storage_join.getRightSampleBlock().getColumnsWithTypeAndName(),
ActionsDAG::MatchColumnsMode::Position);
- QueryPlanStepPtr right_project_step =
std::make_unique<ExpressionStep>(right.getCurrentDataStream(),
std::move(right_project));
+ QueryPlanStepPtr right_project_step =
std::make_unique<ExpressionStep>(right.getCurrentHeader(),
std::move(right_project));
right_project_step->setStepDescription("Rename Broadcast Table Name");
steps.emplace_back(right_project_step.get());
right.addStep(std::move(right_project_step));
@@ -176,17 +176,17 @@ void JoinRelParser::renamePlanColumns(DB::QueryPlan &
left, DB::QueryPlan & righ
/// avoid the columns name in the right table be changed in
`addConvertStep`.
/// This could happen in tpc-ds q44.
DB::ColumnsWithTypeAndName new_left_cols;
- const auto & right_header = right.getCurrentDataStream().header;
+ const auto & right_header = right.getCurrentHeader();
auto left_prefix = getUniqueName("left");
- for (const auto & col : left.getCurrentDataStream().header)
+ for (const auto & col : left.getCurrentHeader())
if (right_header.has(col.name))
new_left_cols.emplace_back(col.column, col.type, left_prefix +
col.name);
else
new_left_cols.emplace_back(col.column, col.type, col.name);
ActionsDAG left_project = ActionsDAG::makeConvertingActions(
- left.getCurrentDataStream().header.getColumnsWithTypeAndName(),
new_left_cols, ActionsDAG::MatchColumnsMode::Position);
+ left.getCurrentHeader().getColumnsWithTypeAndName(), new_left_cols,
ActionsDAG::MatchColumnsMode::Position);
- QueryPlanStepPtr left_project_step =
std::make_unique<ExpressionStep>(left.getCurrentDataStream(),
std::move(left_project));
+ QueryPlanStepPtr left_project_step =
std::make_unique<ExpressionStep>(left.getCurrentHeader(),
std::move(left_project));
left_project_step->setStepDescription("Rename Left Table Name for
broadcast join");
steps.emplace_back(left_project_step.get());
left.addStep(std::move(left_project_step));
@@ -204,14 +204,14 @@ DB::QueryPlanPtr JoinRelParser::parseJoin(const
substrait::JoinRel & join, DB::Q
renamePlanColumns(*left, *right, *storage_join);
auto table_join = createDefaultTableJoin(join.type(),
join_opt_info.is_existence_join, context);
- DB::Block right_header_before_convert_step =
right->getCurrentDataStream().header;
+ DB::Block right_header_before_convert_step = right->getCurrentHeader();
addConvertStep(*table_join, *left, *right);
// Add a check to find error easily.
if (storage_join)
{
bool is_col_names_changed = false;
- const auto & current_right_header =
right->getCurrentDataStream().header;
+ const auto & current_right_header = right->getCurrentHeader();
if (right_header_before_convert_step.columns() !=
current_right_header.columns())
is_col_names_changed = true;
if (!is_col_names_changed)
@@ -230,20 +230,20 @@ DB::QueryPlanPtr JoinRelParser::parseJoin(const
substrait::JoinRel & join, DB::Q
throw DB::Exception(
DB::ErrorCodes::LOGICAL_ERROR,
"For broadcast join, we must not change the columns name in
the right table.\nleft header:{},\nright header: {} -> {}",
- left->getCurrentDataStream().header.dumpStructure(),
+ left->getCurrentHeader().dumpStructure(),
right_header_before_convert_step.dumpStructure(),
- right->getCurrentDataStream().header.dumpStructure());
+ right->getCurrentHeader().dumpStructure());
}
}
Names after_join_names;
- auto left_names = left->getCurrentDataStream().header.getNames();
+ auto left_names = left->getCurrentHeader().getNames();
after_join_names.insert(after_join_names.end(), left_names.begin(),
left_names.end());
auto right_name = table_join->columnsFromJoinedTable().getNames();
after_join_names.insert(after_join_names.end(), right_name.begin(),
right_name.end());
- auto left_header = left->getCurrentDataStream().header;
- auto right_header = right->getCurrentDataStream().header;
+ auto left_header = left->getCurrentHeader();
+ auto right_header = right->getCurrentHeader();
QueryPlanPtr query_plan;
@@ -261,12 +261,12 @@ DB::QueryPlanPtr JoinRelParser::parseJoin(const
substrait::JoinRel & join, DB::Q
if (storage_join->has_null_key_value)
{
// if there is a null key value on the build side, it will
return the empty result
- auto empty_step =
std::make_unique<EarlyStopStep>(left->getCurrentDataStream());
+ auto empty_step =
std::make_unique<EarlyStopStep>(left->getCurrentHeader());
left->addStep(std::move(empty_step));
}
else if (!storage_join->is_empty_hash_table)
{
- auto input_header = left->getCurrentDataStream().header;
+ auto input_header = left->getCurrentHeader();
DB::ActionsDAG
filter_is_not_null_dag{input_header.getColumnsWithTypeAndName()};
// when is_null_aware_anti_join is true, there is only one
join key
const auto * key_field =
filter_is_not_null_dag.getInputs()[join.expression()
@@ -284,7 +284,7 @@ DB::QueryPlanPtr JoinRelParser::parseJoin(const
substrait::JoinRel & join, DB::Q
const auto * cond_node =
buildFunctionNode(filter_is_not_null_dag, "isNotNull", {result_node});
filter_is_not_null_dag.addOrReplaceInOutputs(*cond_node);
auto filter_step = std::make_unique<FilterStep>(
- left->getCurrentDataStream(),
std::move(filter_is_not_null_dag), cond_node->result_name, true);
+ left->getCurrentHeader(),
std::move(filter_is_not_null_dag), cond_node->result_name, true);
left->addStep(std::move(filter_step));
}
// other case: is_empty_hash_table, don't need to handle
@@ -292,7 +292,7 @@ DB::QueryPlanPtr JoinRelParser::parseJoin(const
substrait::JoinRel & join, DB::Q
applyJoinFilter(*table_join, join, *left, *right, true);
auto broadcast_hash_join = storage_join->getJoinLocked(table_join,
context);
- QueryPlanStepPtr join_step =
std::make_unique<FilledJoinStep>(left->getCurrentDataStream(),
broadcast_hash_join, 8192);
+ QueryPlanStepPtr join_step =
std::make_unique<FilledJoinStep>(left->getCurrentHeader(), broadcast_hash_join,
8192);
join_step->setStepDescription("STORAGE_JOIN");
steps.emplace_back(join_step.get());
@@ -311,10 +311,10 @@ DB::QueryPlanPtr JoinRelParser::parseJoin(const
substrait::JoinRel & join, DB::Q
if (need_post_filter && table_join->kind() != DB::JoinKind::Inner)
throw DB::Exception(DB::ErrorCodes::LOGICAL_ERROR, "Sort merge
join doesn't support mixed join conditions, except inner join.");
- JoinPtr smj_join = std::make_shared<FullSortingMergeJoin>(table_join,
right->getCurrentDataStream().header.cloneEmpty(), -1);
+ JoinPtr smj_join = std::make_shared<FullSortingMergeJoin>(table_join,
right->getCurrentHeader().cloneEmpty(), -1);
MultiEnum<DB::JoinAlgorithm> join_algorithm =
context->getSettingsRef()[Setting::join_algorithm];
QueryPlanStepPtr join_step
- = std::make_unique<DB::JoinStep>(left->getCurrentDataStream(),
right->getCurrentDataStream(), smj_join, 8192, 1, false);
+ = std::make_unique<DB::JoinStep>(left->getCurrentHeader(),
right->getCurrentHeader(), smj_join, 8192, 1, false);
join_step->setStepDescription("SORT_MERGE_JOIN");
steps.emplace_back(join_step.get());
@@ -362,7 +362,7 @@ DB::QueryPlanPtr JoinRelParser::parseJoin(const
substrait::JoinRel & join, DB::Q
/// we mark the flag 0, otherwise mark it 1.
void JoinRelParser::existenceJoinPostProject(DB::QueryPlan & plan, const
DB::Names & left_input_cols)
{
- DB::ActionsDAG
actions_dag{plan.getCurrentDataStream().header.getColumnsWithTypeAndName()};
+ DB::ActionsDAG
actions_dag{plan.getCurrentHeader().getColumnsWithTypeAndName()};
const auto * right_col_node = actions_dag.getInputs().back();
auto function_builder = DB::FunctionFactory::instance().get("isNotNull",
getContext());
const auto * not_null_node = &actions_dag.addFunction(function_builder,
{right_col_node}, right_col_node->result_name);
@@ -370,7 +370,7 @@ void JoinRelParser::existenceJoinPostProject(DB::QueryPlan
& plan, const DB::Nam
DB::Names required_cols = left_input_cols;
required_cols.emplace_back(not_null_node->result_name);
actions_dag.removeUnusedActions(required_cols);
- auto project_step =
std::make_unique<DB::ExpressionStep>(plan.getCurrentDataStream(),
std::move(actions_dag));
+ auto project_step =
std::make_unique<DB::ExpressionStep>(plan.getCurrentHeader(),
std::move(actions_dag));
project_step->setStepDescription("ExistenceJoin Post Project");
steps.emplace_back(project_step.get());
plan.addStep(std::move(project_step));
@@ -380,24 +380,23 @@ void JoinRelParser::addConvertStep(TableJoin &
table_join, DB::QueryPlan & left,
{
/// If the columns name in right table is duplicated with left table, we
need to rename the right table's columns.
NameSet left_columns_set;
- for (const auto & col : left.getCurrentDataStream().header.getNames())
+ for (const auto & col : left.getCurrentHeader().getNames())
left_columns_set.emplace(col);
- table_join.setColumnsFromJoinedTable(
- right.getCurrentDataStream().header.getNamesAndTypesList(),
left_columns_set, getUniqueName("right") + ".");
+
table_join.setColumnsFromJoinedTable(right.getCurrentHeader().getNamesAndTypesList(),
left_columns_set, getUniqueName("right") + ".");
// fix right table key duplicate
NamesWithAliases right_table_alias;
for (size_t idx = 0; idx < table_join.columnsFromJoinedTable().size();
idx++)
{
- auto origin_name =
right.getCurrentDataStream().header.getByPosition(idx).name;
+ auto origin_name = right.getCurrentHeader().getByPosition(idx).name;
auto dedup_name =
table_join.columnsFromJoinedTable().getNames().at(idx);
if (origin_name != dedup_name)
right_table_alias.emplace_back(NameWithAlias(origin_name,
dedup_name));
}
if (!right_table_alias.empty())
{
- ActionsDAG
rename_dag{right.getCurrentDataStream().header.getNamesAndTypesList()};
- auto original_right_columns = right.getCurrentDataStream().header;
+ ActionsDAG rename_dag{right.getCurrentHeader().getNamesAndTypesList()};
+ auto original_right_columns = right.getCurrentHeader();
for (const auto & column_alias : right_table_alias)
{
if (original_right_columns.has(column_alias.first))
@@ -408,7 +407,7 @@ void JoinRelParser::addConvertStep(TableJoin & table_join,
DB::QueryPlan & left,
}
}
- QueryPlanStepPtr project_step =
std::make_unique<ExpressionStep>(right.getCurrentDataStream(),
std::move(rename_dag));
+ QueryPlanStepPtr project_step =
std::make_unique<ExpressionStep>(right.getCurrentHeader(),
std::move(rename_dag));
project_step->setStepDescription("Right Table Rename");
steps.emplace_back(project_step.get());
right.addStep(std::move(project_step));
@@ -419,11 +418,11 @@ void JoinRelParser::addConvertStep(TableJoin &
table_join, DB::QueryPlan & left,
std::optional<ActionsDAG> left_convert_actions;
std::optional<ActionsDAG> right_convert_actions;
std::tie(left_convert_actions, right_convert_actions) =
table_join.createConvertingActions(
- left.getCurrentDataStream().header.getColumnsWithTypeAndName(),
right.getCurrentDataStream().header.getColumnsWithTypeAndName());
+ left.getCurrentHeader().getColumnsWithTypeAndName(),
right.getCurrentHeader().getColumnsWithTypeAndName());
if (right_convert_actions)
{
- auto converting_step =
std::make_unique<ExpressionStep>(right.getCurrentDataStream(),
std::move(*right_convert_actions));
+ auto converting_step =
std::make_unique<ExpressionStep>(right.getCurrentHeader(),
std::move(*right_convert_actions));
converting_step->setStepDescription("Convert joined columns");
steps.emplace_back(converting_step.get());
right.addStep(std::move(converting_step));
@@ -431,7 +430,7 @@ void JoinRelParser::addConvertStep(TableJoin & table_join,
DB::QueryPlan & left,
if (left_convert_actions)
{
- auto converting_step =
std::make_unique<ExpressionStep>(left.getCurrentDataStream(),
std::move(*left_convert_actions));
+ auto converting_step =
std::make_unique<ExpressionStep>(left.getCurrentHeader(),
std::move(*left_convert_actions));
converting_step->setStepDescription("Convert joined columns");
steps.emplace_back(converting_step.get());
left.addStep(std::move(converting_step));
@@ -510,8 +509,8 @@ bool JoinRelParser::applyJoinFilter(
return true;
const auto & expr = join_rel.post_join_filter();
- const auto & left_header = left.getCurrentDataStream().header;
- const auto & right_header = right.getCurrentDataStream().header;
+ const auto & left_header = left.getCurrentHeader();
+ const auto & right_header = right.getCurrentHeader();
ColumnsWithTypeAndName mixed_columns;
std::unordered_set<String> added_column_name;
for (const auto & col : left_header.getColumnsWithTypeAndName())
@@ -557,7 +556,7 @@ bool JoinRelParser::applyJoinFilter(
input_exprs.push_back(expr);
auto actions_dag = expressionsToActionsDAG(input_exprs,
left_header);
table_join.getClauses().back().analyzer_left_filter_condition_column_name =
actions_dag.getOutputs().back()->result_name;
- QueryPlanStepPtr before_join_step =
std::make_unique<ExpressionStep>(left.getCurrentDataStream(),
std::move(actions_dag));
+ QueryPlanStepPtr before_join_step =
std::make_unique<ExpressionStep>(left.getCurrentHeader(),
std::move(actions_dag));
before_join_step->setStepDescription("Before JOIN LEFT");
steps.emplace_back(before_join_step.get());
left.addStep(std::move(before_join_step));
@@ -576,7 +575,7 @@ bool JoinRelParser::applyJoinFilter(
actions_dag.removeUnusedActions();
table_join.getClauses().back().analyzer_right_filter_condition_column_name =
actions_dag.getOutputs().back()->result_name;
- QueryPlanStepPtr before_join_step =
std::make_unique<ExpressionStep>(right.getCurrentDataStream(),
std::move(actions_dag));
+ QueryPlanStepPtr before_join_step =
std::make_unique<ExpressionStep>(right.getCurrentHeader(),
std::move(actions_dag));
before_join_step->setStepDescription("Before JOIN RIGHT");
steps.emplace_back(before_join_step.get());
right.addStep(std::move(before_join_step));
@@ -602,7 +601,7 @@ bool JoinRelParser::applyJoinFilter(
void JoinRelParser::addPostFilter(DB::QueryPlan & query_plan, const
substrait::JoinRel & join)
{
std::string filter_name;
- ActionsDAG
actions_dag{query_plan.getCurrentDataStream().header.getColumnsWithTypeAndName()};
+ ActionsDAG
actions_dag{query_plan.getCurrentHeader().getColumnsWithTypeAndName()};
if (!join.post_join_filter().has_scalar_function())
{
// It may be singular_or_list
@@ -614,7 +613,7 @@ void JoinRelParser::addPostFilter(DB::QueryPlan &
query_plan, const substrait::J
const auto * func_node =
expression_parser->parseFunction(join.post_join_filter().scalar_function(),
actions_dag, true);
filter_name = func_node->result_name;
}
- auto filter_step =
std::make_unique<FilterStep>(query_plan.getCurrentDataStream(),
std::move(actions_dag), filter_name, true);
+ auto filter_step =
std::make_unique<FilterStep>(query_plan.getCurrentHeader(),
std::move(actions_dag), filter_name, true);
filter_step->setStepDescription("Post Join Filter");
steps.emplace_back(filter_step.get());
query_plan.addStep(std::move(filter_step));
@@ -772,9 +771,9 @@ DB::QueryPlanPtr JoinRelParser::buildMultiOnClauseHashJoin(
LOG_INFO(getLogger("JoinRelParser"), "multi join on clauses:\n{}",
DB::TableJoin::formatClauses(table_join->getClauses()));
- JoinPtr hash_join = std::make_shared<HashJoin>(table_join,
right_plan->getCurrentDataStream().header);
+ JoinPtr hash_join = std::make_shared<HashJoin>(table_join,
right_plan->getCurrentHeader());
QueryPlanStepPtr join_step
- = std::make_unique<DB::JoinStep>(left_plan->getCurrentDataStream(),
right_plan->getCurrentDataStream(), hash_join, 8192, 1, false);
+ = std::make_unique<DB::JoinStep>(left_plan->getCurrentHeader(),
right_plan->getCurrentHeader(), hash_join, 8192, 1, false);
join_step->setStepDescription("Multi join on clause hash join");
steps.emplace_back(join_step.get());
std::vector<QueryPlanPtr> plans;
@@ -801,18 +800,14 @@ DB::QueryPlanPtr
JoinRelParser::buildSingleOnClauseHashJoin(
if (join_algorithm.isSet(DB::JoinAlgorithm::GRACE_HASH))
{
hash_join = std::make_shared<GraceHashJoin>(
- context,
- table_join,
- left_plan->getCurrentDataStream().header,
- right_plan->getCurrentDataStream().header,
- context->getTempDataOnDisk());
+ context, table_join, left_plan->getCurrentHeader(),
right_plan->getCurrentHeader(), context->getTempDataOnDisk());
}
else
{
- hash_join = std::make_shared<HashJoin>(table_join,
right_plan->getCurrentDataStream().header.cloneEmpty());
+ hash_join = std::make_shared<HashJoin>(table_join,
right_plan->getCurrentHeader().cloneEmpty());
}
QueryPlanStepPtr join_step
- = std::make_unique<DB::JoinStep>(left_plan->getCurrentDataStream(),
right_plan->getCurrentDataStream(), hash_join, 8192, 1, false);
+ = std::make_unique<DB::JoinStep>(left_plan->getCurrentHeader(),
right_plan->getCurrentHeader(), hash_join, 8192, 1, false);
join_step->setStepDescription("HASH_JOIN");
steps.emplace_back(join_step.get());
diff --git a/cpp-ch/local-engine/Parser/RelParsers/MergeTreeRelParser.cpp
b/cpp-ch/local-engine/Parser/RelParsers/MergeTreeRelParser.cpp
index 2155c385b0..7947fc2536 100644
--- a/cpp-ch/local-engine/Parser/RelParsers/MergeTreeRelParser.cpp
+++ b/cpp-ch/local-engine/Parser/RelParsers/MergeTreeRelParser.cpp
@@ -153,7 +153,7 @@ DB::QueryPlanPtr MergeTreeRelParser::parseReadRel(
query_plan->addStep(std::move(read_step));
if (!non_nullable_columns.empty())
{
- auto input_header = query_plan->getCurrentDataStream().header;
+ auto input_header = query_plan->getCurrentHeader();
std::erase_if(non_nullable_columns, [input_header](auto item) -> bool
{ return !input_header.has(item); });
auto * remove_null_step =
PlanUtil::addRemoveNullableStep(parser_context->queryContext(), *query_plan,
non_nullable_columns);
if (remove_null_step)
diff --git a/cpp-ch/local-engine/Parser/RelParsers/ProjectRelParser.cpp
b/cpp-ch/local-engine/Parser/RelParsers/ProjectRelParser.cpp
index 4a6bd2f488..20af6f83fc 100644
--- a/cpp-ch/local-engine/Parser/RelParsers/ProjectRelParser.cpp
+++ b/cpp-ch/local-engine/Parser/RelParsers/ProjectRelParser.cpp
@@ -63,13 +63,13 @@ ProjectRelParser::parseProject(DB::QueryPlanPtr query_plan,
const substrait::Rel
if (project_rel.expressions_size())
{
std::vector<substrait::Expression> expressions;
- auto header = query_plan->getCurrentDataStream().header;
+ auto header = query_plan->getCurrentHeader();
for (int i = 0; i < project_rel.expressions_size(); ++i)
{
expressions.emplace_back(project_rel.expressions(i));
}
auto actions_dag = expressionsToActionsDAG(expressions, header);
- auto expression_step =
std::make_unique<ExpressionStep>(query_plan->getCurrentDataStream(),
std::move(actions_dag));
+ auto expression_step =
std::make_unique<ExpressionStep>(query_plan->getCurrentHeader(),
std::move(actions_dag));
expression_step->setStepDescription("Project");
steps.emplace_back(expression_step.get());
query_plan->addStep(std::move(expression_step));
@@ -77,7 +77,7 @@ ProjectRelParser::parseProject(DB::QueryPlanPtr query_plan,
const substrait::Rel
}
else
{
- auto empty_project_step =
std::make_unique<EmptyProjectStep>(query_plan->getCurrentDataStream());
+ auto empty_project_step =
std::make_unique<EmptyProjectStep>(query_plan->getCurrentHeader());
empty_project_step->setStepDescription("EmptyProject");
steps.emplace_back(empty_project_step.get());
query_plan->addStep(std::move(empty_project_step));
@@ -132,14 +132,14 @@ DB::QueryPlanPtr
ProjectRelParser::parseReplicateRows(DB::QueryPlanPtr query_pla
{
expressions.emplace_back(generate_rel.generator().scalar_function().arguments(i).value());
}
- auto header = query_plan->getCurrentDataStream().header;
+ auto header = query_plan->getCurrentHeader();
auto actions_dag = expressionsToActionsDAG(expressions, header);
- auto before_replicate_rows =
std::make_unique<DB::ExpressionStep>(query_plan->getCurrentDataStream(),
std::move(actions_dag));
+ auto before_replicate_rows =
std::make_unique<DB::ExpressionStep>(query_plan->getCurrentHeader(),
std::move(actions_dag));
before_replicate_rows->setStepDescription("Before ReplicateRows");
steps.emplace_back(before_replicate_rows.get());
query_plan->addStep(std::move(before_replicate_rows));
- auto replicate_rows_step =
std::make_unique<ReplicateRowsStep>(query_plan->getCurrentDataStream());
+ auto replicate_rows_step =
std::make_unique<ReplicateRowsStep>(query_plan->getCurrentHeader());
replicate_rows_step->setStepDescription("ReplicateRows");
steps.emplace_back(replicate_rows_step.get());
query_plan->addStep(std::move(replicate_rows_step));
@@ -161,13 +161,13 @@ ProjectRelParser::parseGenerate(DB::QueryPlanPtr
query_plan, const substrait::Re
}
expressions.emplace_back(generate_rel.generator());
- auto header = query_plan->getCurrentDataStream().header;
+ auto header = query_plan->getCurrentHeader();
auto actions_dag = expressionsToActionsDAG(expressions, header);
if (!findArrayJoinNode(actions_dag))
{
/// If generator in generate rel is not explode/posexplode, e.g.
json_tuple
- auto expression_step =
std::make_unique<ExpressionStep>(query_plan->getCurrentDataStream(),
std::move(actions_dag));
+ auto expression_step =
std::make_unique<ExpressionStep>(query_plan->getCurrentHeader(),
std::move(actions_dag));
expression_step->setStepDescription("Generate");
steps.emplace_back(expression_step.get());
query_plan->addStep(std::move(expression_step));
@@ -198,7 +198,7 @@ ProjectRelParser::parseGenerate(DB::QueryPlanPtr
query_plan, const substrait::Re
if (!ignore_actions_dag(splitted_actions_dags.before_array_join))
{
auto step_before_array_join
- =
std::make_unique<ExpressionStep>(query_plan->getCurrentDataStream(),
std::move(splitted_actions_dags.before_array_join));
+ =
std::make_unique<ExpressionStep>(query_plan->getCurrentHeader(),
std::move(splitted_actions_dags.before_array_join));
step_before_array_join->setStepDescription("Pre-projection In
Generate");
steps.emplace_back(step_before_array_join.get());
query_plan->addStep(std::move(step_before_array_join));
@@ -211,7 +211,7 @@ ProjectRelParser::parseGenerate(DB::QueryPlanPtr
query_plan, const substrait::Re
array_join.columns = std::move(array_joined_columns);
array_join.is_left = generate_rel.outer();
auto array_join_step = std::make_unique<ArrayJoinStep>(
- query_plan->getCurrentDataStream(), std::move(array_join), false,
getContext()->getSettingsRef()[Setting::max_block_size]);
+ query_plan->getCurrentHeader(), std::move(array_join), false,
getContext()->getSettingsRef()[Setting::max_block_size]);
array_join_step->setStepDescription("ARRAY JOIN In Generate");
steps.emplace_back(array_join_step.get());
query_plan->addStep(std::move(array_join_step));
@@ -221,7 +221,7 @@ ProjectRelParser::parseGenerate(DB::QueryPlanPtr
query_plan, const substrait::Re
if (!ignore_actions_dag(splitted_actions_dags.after_array_join))
{
auto step_after_array_join
- =
std::make_unique<ExpressionStep>(query_plan->getCurrentDataStream(),
std::move(splitted_actions_dags.after_array_join));
+ =
std::make_unique<ExpressionStep>(query_plan->getCurrentHeader(),
std::move(splitted_actions_dags.after_array_join));
step_after_array_join->setStepDescription("Post-projection In
Generate");
steps.emplace_back(step_after_array_join.get());
query_plan->addStep(std::move(step_after_array_join));
diff --git a/cpp-ch/local-engine/Parser/RelParsers/ReadRelParser.cpp
b/cpp-ch/local-engine/Parser/RelParsers/ReadRelParser.cpp
index c41c07b380..75e6e14c4a 100644
--- a/cpp-ch/local-engine/Parser/RelParsers/ReadRelParser.cpp
+++ b/cpp-ch/local-engine/Parser/RelParsers/ReadRelParser.cpp
@@ -64,7 +64,7 @@ DB::QueryPlanPtr ReadRelParser::parse(DB::QueryPlanPtr
query_plan, const substra
if (getContext()->getSettingsRef()[Setting::max_threads] > 1)
{
- auto buffer_step =
std::make_unique<BlocksBufferPoolStep>(query_plan->getCurrentDataStream());
+ auto buffer_step =
std::make_unique<BlocksBufferPoolStep>(query_plan->getCurrentHeader());
steps.emplace_back(buffer_step.get());
query_plan->addStep(std::move(buffer_step));
}
diff --git a/cpp-ch/local-engine/Parser/RelParsers/SortRelParser.cpp
b/cpp-ch/local-engine/Parser/RelParsers/SortRelParser.cpp
index 552904956a..1ed4f2565d 100644
--- a/cpp-ch/local-engine/Parser/RelParsers/SortRelParser.cpp
+++ b/cpp-ch/local-engine/Parser/RelParsers/SortRelParser.cpp
@@ -41,16 +41,12 @@ SortRelParser::parse(DB::QueryPlanPtr query_plan, const
substrait::Rel & rel, st
{
size_t limit = parseLimit(rel_stack_);
const auto & sort_rel = rel.sort();
- auto sort_descr = parseSortDescription(sort_rel.sorts(),
query_plan->getCurrentDataStream().header);
+ auto sort_descr = parseSortDescription(sort_rel.sorts(),
query_plan->getCurrentHeader());
SortingStep::Settings settings(*getContext());
auto config = MemoryConfig::loadFromContext(getContext());
double spill_mem_ratio = config.spill_mem_ratio;
- settings.worth_external_sort = [spill_mem_ratio]() -> bool
- {
- return currentThreadGroupMemoryUsageRatio() > spill_mem_ratio;
- };
- auto sorting_step = std::make_unique<DB::SortingStep>(
- query_plan->getCurrentDataStream(), sort_descr, limit, settings);
+ settings.worth_external_sort = [spill_mem_ratio]() -> bool { return
currentThreadGroupMemoryUsageRatio() > spill_mem_ratio; };
+ auto sorting_step =
std::make_unique<DB::SortingStep>(query_plan->getCurrentHeader(), sort_descr,
limit, settings);
sorting_step->setStepDescription("Sorting step");
steps.emplace_back(sorting_step.get());
query_plan->addStep(std::move(sorting_step));
diff --git
a/cpp-ch/local-engine/Parser/RelParsers/WindowGroupLimitRelParser.cpp
b/cpp-ch/local-engine/Parser/RelParsers/WindowGroupLimitRelParser.cpp
index 44d7053c7e..e82d68c1d1 100644
--- a/cpp-ch/local-engine/Parser/RelParsers/WindowGroupLimitRelParser.cpp
+++ b/cpp-ch/local-engine/Parser/RelParsers/WindowGroupLimitRelParser.cpp
@@ -49,7 +49,7 @@ WindowGroupLimitRelParser::parse(DB::QueryPlanPtr
current_plan_, const substrait
size_t limit = static_cast<size_t>(win_rel_def.limit());
auto window_group_limit_step = std::make_unique<WindowGroupLimitStep>(
- current_plan->getCurrentDataStream(), window_function_name,
partition_fields, sort_fields, limit);
+ current_plan->getCurrentHeader(), window_function_name,
partition_fields, sort_fields, limit);
window_group_limit_step->setStepDescription("Window group limit");
steps.emplace_back(window_group_limit_step.get());
current_plan->addStep(std::move(window_group_limit_step));
@@ -62,20 +62,12 @@ WindowGroupLimitRelParser::parsePartitoinFields(const
google::protobuf::Repeated
{
std::vector<size_t> fields;
for (const auto & expr : expressions)
- {
if (expr.has_selection())
- {
fields.push_back(static_cast<size_t>(expr.selection().direct_reference().struct_field().field()));
- }
else if (expr.has_literal())
- {
continue;
- }
else
- {
throw DB::Exception(DB::ErrorCodes::BAD_ARGUMENTS, "Unknow
expression: {}", expr.DebugString());
- }
- }
return fields;
}
@@ -83,20 +75,12 @@ std::vector<size_t>
WindowGroupLimitRelParser::parseSortFields(const google::pro
{
std::vector<size_t> fields;
for (const auto sort_field : sort_fields)
- {
if (sort_field.expr().has_literal())
- {
continue;
- }
else if (sort_field.expr().has_selection())
- {
fields.push_back(static_cast<size_t>(sort_field.expr().selection().direct_reference().struct_field().field()));
- }
else
- {
throw DB::Exception(DB::ErrorCodes::BAD_ARGUMENTS, "Unknown
expression: {}", sort_field.expr().DebugString());
- }
- }
return fields;
}
diff --git a/cpp-ch/local-engine/Parser/RelParsers/WindowRelParser.cpp
b/cpp-ch/local-engine/Parser/RelParsers/WindowRelParser.cpp
index 00d1c5c3f7..7b5b0147ab 100644
--- a/cpp-ch/local-engine/Parser/RelParsers/WindowRelParser.cpp
+++ b/cpp-ch/local-engine/Parser/RelParsers/WindowRelParser.cpp
@@ -59,7 +59,7 @@ WindowRelParser::parse(DB::QueryPlanPtr current_plan_, const
substrait::Rel & re
{
const auto & win_rel_pb = rel.window();
current_plan = std::move(current_plan_);
- input_header = current_plan->getCurrentDataStream().header;
+ input_header = current_plan->getCurrentHeader();
// The output header is : original columns ++ window columns
output_header = input_header;
for (const auto & measure : win_rel_pb.measures())
@@ -82,7 +82,7 @@ WindowRelParser::parse(DB::QueryPlanPtr current_plan_, const
substrait::Rel & re
{
auto & win = it.second;
- auto window_step =
std::make_unique<DB::WindowStep>(current_plan->getCurrentDataStream(), win,
win.window_functions, false);
+ auto window_step =
std::make_unique<DB::WindowStep>(current_plan->getCurrentHeader(), win,
win.window_functions, false);
window_step->setStepDescription("Window step for window '" +
win.window_name + "'");
steps.emplace_back(window_step.get());
current_plan->addStep(std::move(window_step));
@@ -98,7 +98,7 @@ DB::WindowDescription
WindowRelParser::parseWindowDescription(const WindowInfo &
DB::WindowDescription win_descr;
win_descr.frame = parseWindowFrame(win_info);
win_descr.partition_by = parsePartitionBy(win_info.partition_exprs);
- win_descr.order_by =
SortRelParser::parseSortDescription(win_info.sort_fields,
current_plan->getCurrentDataStream().header);
+ win_descr.order_by =
SortRelParser::parseSortDescription(win_info.sort_fields,
current_plan->getCurrentHeader());
win_descr.full_sort_description = win_descr.partition_by;
win_descr.full_sort_description.insert(win_descr.full_sort_description.end(),
win_descr.order_by.begin(), win_descr.order_by.end());
@@ -254,7 +254,7 @@ void WindowRelParser::parseBoundType(
DB::SortDescription WindowRelParser::parsePartitionBy(const
google::protobuf::RepeatedPtrField<substrait::Expression> & expressions)
{
- DB::Block header = current_plan->getCurrentDataStream().header;
+ DB::Block header = current_plan->getCurrentHeader();
DB::SortDescription sort_descr;
for (const auto & expr : expressions)
{
@@ -316,7 +316,7 @@ void WindowRelParser::initWindowsInfos(const
substrait::WindowRel & win_rel)
void WindowRelParser::tryAddProjectionBeforeWindow()
{
- auto header = current_plan->getCurrentDataStream().header;
+ auto header = current_plan->getCurrentHeader();
ActionsDAG actions_dag{header.getColumnsWithTypeAndName()};
auto dag_footprint = actions_dag.dumpDAG();
@@ -335,7 +335,7 @@ void WindowRelParser::tryAddProjectionBeforeWindow()
if (actions_dag.dumpDAG() != dag_footprint)
{
- auto project_step =
std::make_unique<ExpressionStep>(current_plan->getCurrentDataStream(),
std::move(actions_dag));
+ auto project_step =
std::make_unique<ExpressionStep>(current_plan->getCurrentHeader(),
std::move(actions_dag));
project_step->setStepDescription("Add projections before window");
steps.emplace_back(project_step.get());
current_plan->addStep(std::move(project_step));
@@ -345,7 +345,7 @@ void WindowRelParser::tryAddProjectionBeforeWindow()
void WindowRelParser::tryAddProjectionAfterWindow()
{
// The final result header is : original header ++ [window aggregate
columns]
- auto header = current_plan->getCurrentDataStream().header;
+ auto header = current_plan->getCurrentHeader();
ActionsDAG actions_dag{header.getColumnsWithTypeAndName()};
auto dag_footprint = actions_dag.dumpDAG();
@@ -358,20 +358,20 @@ void WindowRelParser::tryAddProjectionAfterWindow()
if (actions_dag.dumpDAG() != dag_footprint)
{
- auto project_step =
std::make_unique<ExpressionStep>(current_plan->getCurrentDataStream(),
std::move(actions_dag));
+ auto project_step =
std::make_unique<ExpressionStep>(current_plan->getCurrentHeader(),
std::move(actions_dag));
project_step->setStepDescription("Add projections for window result");
steps.emplace_back(project_step.get());
current_plan->addStep(std::move(project_step));
}
// This projeciton will remove the const columns from the window function
arguments
- auto current_header = current_plan->getCurrentDataStream().header;
+ auto current_header = current_plan->getCurrentHeader();
if (!DB::blocksHaveEqualStructure(output_header, current_header))
{
ActionsDAG convert_action = ActionsDAG::makeConvertingActions(
current_header.getColumnsWithTypeAndName(),
output_header.getColumnsWithTypeAndName(),
DB::ActionsDAG::MatchColumnsMode::Name);
QueryPlanStepPtr convert_step
- =
std::make_unique<DB::ExpressionStep>(current_plan->getCurrentDataStream(),
std::move(convert_action));
+ =
std::make_unique<DB::ExpressionStep>(current_plan->getCurrentHeader(),
std::move(convert_action));
convert_step->setStepDescription("Convert window Output");
steps.emplace_back(convert_step.get());
current_plan->addStep(std::move(convert_step));
diff --git a/cpp-ch/local-engine/Parser/SerializedPlanParser.cpp
b/cpp-ch/local-engine/Parser/SerializedPlanParser.cpp
index 8b305488b6..32ead10708 100644
--- a/cpp-ch/local-engine/Parser/SerializedPlanParser.cpp
+++ b/cpp-ch/local-engine/Parser/SerializedPlanParser.cpp
@@ -119,9 +119,9 @@ void adjustOutput(const DB::QueryPlanPtr & query_plan,
const substrait::PlanRel
{
if (root_rel.root().names_size())
{
- ActionsDAG
actions_dag{blockToNameAndTypeList(query_plan->getCurrentDataStream().header)};
+ ActionsDAG
actions_dag{blockToNameAndTypeList(query_plan->getCurrentHeader())};
NamesWithAliases aliases;
- auto cols =
query_plan->getCurrentDataStream().header.getNamesAndTypesList();
+ auto cols = query_plan->getCurrentHeader().getNamesAndTypesList();
if (cols.getNames().size() !=
static_cast<size_t>(root_rel.root().names_size()))
throw Exception(
ErrorCodes::LOGICAL_ERROR,
@@ -131,7 +131,7 @@ void adjustOutput(const DB::QueryPlanPtr & query_plan,
const substrait::PlanRel
for (int i = 0; i < static_cast<int>(cols.getNames().size()); i++)
aliases.emplace_back(NameWithAlias(cols.getNames()[i],
root_rel.root().names(i)));
actions_dag.project(aliases);
- auto expression_step =
std::make_unique<ExpressionStep>(query_plan->getCurrentDataStream(),
std::move(actions_dag));
+ auto expression_step =
std::make_unique<ExpressionStep>(query_plan->getCurrentHeader(),
std::move(actions_dag));
expression_step->setStepDescription("Rename Output");
query_plan->addStep(std::move(expression_step));
}
@@ -140,10 +140,18 @@ void adjustOutput(const DB::QueryPlanPtr & query_plan,
const substrait::PlanRel
const auto & output_schema = root_rel.root().output_schema();
if (output_schema.types_size())
{
- auto original_header = query_plan->getCurrentDataStream().header;
+ auto original_header = query_plan->getCurrentHeader();
const auto & original_cols =
original_header.getColumnsWithTypeAndName();
if (static_cast<size_t>(output_schema.types_size()) !=
original_cols.size())
- throw Exception(ErrorCodes::LOGICAL_ERROR, "Mismatch output
schema");
+ {
+ throw Exception(
+ ErrorCodes::LOGICAL_ERROR,
+ "Mismatch output schema. plan column size {} [header: '{}'],
subtrait plan size {}[schema: {}].",
+ original_cols.size(),
+ original_header.dumpStructure(),
+ output_schema.types_size(),
+ dumpMessage(output_schema));
+ }
bool need_final_project = false;
ColumnsWithTypeAndName final_cols;
for (int i = 0; i < output_schema.types_size(); ++i)
@@ -175,7 +183,7 @@ void adjustOutput(const DB::QueryPlanPtr & query_plan,
const substrait::PlanRel
{
ActionsDAG final_project =
ActionsDAG::makeConvertingActions(original_cols, final_cols,
ActionsDAG::MatchColumnsMode::Position);
QueryPlanStepPtr final_project_step
- =
std::make_unique<ExpressionStep>(query_plan->getCurrentDataStream(),
std::move(final_project));
+ =
std::make_unique<ExpressionStep>(query_plan->getCurrentHeader(),
std::move(final_project));
final_project_step->setStepDescription("Project for output
schema");
query_plan->addStep(std::move(final_project_step));
}
diff --git a/cpp-ch/local-engine/Parser/SubstraitParserUtils.cpp
b/cpp-ch/local-engine/Parser/SubstraitParserUtils.cpp
index 164cf0392d..c16405eff3 100644
--- a/cpp-ch/local-engine/Parser/SubstraitParserUtils.cpp
+++ b/cpp-ch/local-engine/Parser/SubstraitParserUtils.cpp
@@ -24,18 +24,30 @@ using namespace DB;
namespace local_engine
{
+namespace pb_util = google::protobuf::util;
void logDebugMessage(const google::protobuf::Message & message, const char *
type)
{
if (auto * logger = &Poco::Logger::get("SubstraitPlan"); logger->debug())
{
- namespace pb_util = google::protobuf::util;
pb_util::JsonOptions options;
std::string json;
- if (auto s = pb_util::MessageToJsonString(message, &json, options);
!s.ok())
+ if (auto s = MessageToJsonString(message, &json, options); !s.ok())
throw Exception(ErrorCodes::LOGICAL_ERROR, "Can not convert {} to
Json", type);
LOG_DEBUG(logger, "{}:\n{}", type, json);
}
}
+std::string dumpMessage(const google::protobuf::Message & message)
+{
+ pb_util::JsonOptions options;
+ std::string json;
+ if (auto s = MessageToJsonString(message, &json, options); !s.ok())
+ {
+ if (auto * logger = &Poco::Logger::get("SubstraitPlan"))
+ LOG_ERROR(logger, "Can not convert message to Json");
+ return "";
+ }
+ return json;
+}
std::string toString(const google::protobuf::Any & any)
{
google::protobuf::StringValue sv;
diff --git a/cpp-ch/local-engine/Parser/SubstraitParserUtils.h
b/cpp-ch/local-engine/Parser/SubstraitParserUtils.h
index d93b80cdac..a6c252034c 100644
--- a/cpp-ch/local-engine/Parser/SubstraitParserUtils.h
+++ b/cpp-ch/local-engine/Parser/SubstraitParserUtils.h
@@ -69,5 +69,7 @@ Message BinaryToMessage(const std::string_view binary)
void logDebugMessage(const google::protobuf::Message & message, const char *
type);
+std::string dumpMessage(const google::protobuf::Message & message);
+
std::string toString(const google::protobuf::Any & any);
} // namespace local_engine
diff --git
a/cpp-ch/local-engine/Storages/Parquet/VectorizedParquetRecordReader.cpp
b/cpp-ch/local-engine/Storages/Parquet/VectorizedParquetRecordReader.cpp
index 8c9ac073a3..5208ace665 100644
--- a/cpp-ch/local-engine/Storages/Parquet/VectorizedParquetRecordReader.cpp
+++ b/cpp-ch/local-engine/Storages/Parquet/VectorizedParquetRecordReader.cpp
@@ -196,7 +196,17 @@ bool VectorizedParquetRecordReader::initialize(
/// column pruning
DB::ArrowFieldIndexUtil field_util(
format_settings_.parquet.case_insensitive_column_matching,
format_settings_.parquet.allow_missing_columns);
- const std::vector<Int32> column_indices =
field_util.findRequiredIndices(header, schema);
+ auto index_mapping = field_util.findRequiredIndices(header, schema,
*metadata);
+
+ std::vector<Int32> column_indices;
+ for (const auto & [clickhouse_header_index, parquet_indexes] :
index_mapping)
+ {
+ for (auto parquet_index : parquet_indexes)
+ {
+ column_indices.push_back(parquet_index);
+ }
+ }
+
THROW_ARROW_NOT_OK_OR_ASSIGN(std::vector<int> field_indices,
manifest.GetFieldIndices(column_indices));
/// row groups pruning
@@ -367,7 +377,7 @@ void VectorizedParquetBlockInputFormat::resetParser()
{
IInputFormat::resetParser();
record_reader_.reset();
- }
+}
DB::Chunk VectorizedParquetBlockInputFormat::read()
{
diff --git
a/cpp-ch/local-engine/Storages/SubstraitSource/SubstraitFileSourceStep.cpp
b/cpp-ch/local-engine/Storages/SubstraitSource/SubstraitFileSourceStep.cpp
index 8f7586d290..91e50b4275 100644
--- a/cpp-ch/local-engine/Storages/SubstraitSource/SubstraitFileSourceStep.cpp
+++ b/cpp-ch/local-engine/Storages/SubstraitSource/SubstraitFileSourceStep.cpp
@@ -52,9 +52,7 @@ SubstraitFileStorage dummy_storage{DB::StorageID("dummy_db",
"dummy_table")};
}
SubstraitFileSourceStep::SubstraitFileSourceStep(const DB::ContextPtr &
context_, DB::Pipe pipe_, const String &)
- : SourceStepWithFilter(
- DB::DataStream{.header = pipe_.getHeader()}, {}, {},
dummy_storage.getStorageSnapshot(nullptr, nullptr), context_)
- , pipe(std::move(pipe_))
+ : SourceStepWithFilter(pipe_.getHeader(), {}, {},
dummy_storage.getStorageSnapshot(nullptr, nullptr), context_),
pipe(std::move(pipe_))
{
}
diff --git a/cpp-ch/local-engine/tests/benchmark_local_engine.cpp
b/cpp-ch/local-engine/tests/benchmark_local_engine.cpp
index cee4035fa4..13e74abaee 100644
--- a/cpp-ch/local-engine/tests/benchmark_local_engine.cpp
+++ b/cpp-ch/local-engine/tests/benchmark_local_engine.cpp
@@ -24,8 +24,9 @@
#include <Interpreters/HashJoin/HashJoin.h>
#include <Interpreters/TableJoin.h>
#include <Parser/CHColumnToSparkRow.h>
-#include <Parser/SerializedPlanParser.h>
+#include <Parser/LocalExecutor.h>
#include <Parser/ParserContext.h>
+#include <Parser/SerializedPlanParser.h>
#include <Parser/SparkRowToCHColumn.h>
#include <Parser/SubstraitParserUtils.h>
#include <Parsers/ASTIdentifier.h>
@@ -43,11 +44,10 @@
#include <Common/CHUtil.h>
#include <Common/DebugUtils.h>
#include <Common/PODArray_fwd.h>
+#include <Common/QueryContext.h>
#include <Common/Stopwatch.h>
#include <Common/logger_useful.h>
-#include <Parser/LocalExecutor.h>
#include "testConfig.h"
-#include <Common/QueryContext.h>
#if defined(__SSE2__)
#include <emmintrin.h>
@@ -812,11 +812,11 @@ QueryPlanPtr joinPlan(QueryPlanPtr left, QueryPlanPtr
right, String left_key, St
{
auto join = std::make_shared<TableJoin>(
global_context->getSettingsRef(),
global_context->getGlobalTemporaryVolume(),
global_context->getTempDataOnDisk());
- auto left_columns =
left->getCurrentDataStream().header.getColumnsWithTypeAndName();
- auto right_columns =
right->getCurrentDataStream().header.getColumnsWithTypeAndName();
+ auto left_columns = left->getCurrentHeader().getColumnsWithTypeAndName();
+ auto right_columns = right->getCurrentHeader().getColumnsWithTypeAndName();
join->setKind(JoinKind::Left);
join->setStrictness(JoinStrictness::All);
-
join->setColumnsFromJoinedTable(right->getCurrentDataStream().header.getNamesAndTypesList());
+
join->setColumnsFromJoinedTable(right->getCurrentHeader().getNamesAndTypesList());
join->addDisjunct();
ASTPtr lkey = std::make_shared<ASTIdentifier>(left_key);
ASTPtr rkey = std::make_shared<ASTIdentifier>(right_key);
@@ -824,7 +824,7 @@ QueryPlanPtr joinPlan(QueryPlanPtr left, QueryPlanPtr
right, String left_key, St
for (const auto & column : join->columnsFromJoinedTable())
join->addJoinedColumn(column);
- auto left_keys =
left->getCurrentDataStream().header.getNamesAndTypesList();
+ auto left_keys = left->getCurrentHeader().getNamesAndTypesList();
join->addJoinedColumnsAndCorrectTypes(left_keys, true);
std::optional<ActionsDAG> left_convert_actions;
std::optional<ActionsDAG> right_convert_actions;
@@ -832,21 +832,21 @@ QueryPlanPtr joinPlan(QueryPlanPtr left, QueryPlanPtr
right, String left_key, St
if (right_convert_actions)
{
- auto converting_step =
std::make_unique<ExpressionStep>(right->getCurrentDataStream(),
std::move(*right_convert_actions));
+ auto converting_step =
std::make_unique<ExpressionStep>(right->getCurrentHeader(),
std::move(*right_convert_actions));
converting_step->setStepDescription("Convert joined columns");
right->addStep(std::move(converting_step));
}
if (left_convert_actions)
{
- auto converting_step =
std::make_unique<ExpressionStep>(right->getCurrentDataStream(),
std::move(*right_convert_actions));
+ auto converting_step =
std::make_unique<ExpressionStep>(right->getCurrentHeader(),
std::move(*right_convert_actions));
converting_step->setStepDescription("Convert joined columns");
left->addStep(std::move(converting_step));
}
- auto hash_join = std::make_shared<HashJoin>(join,
right->getCurrentDataStream().header);
+ auto hash_join = std::make_shared<HashJoin>(join,
right->getCurrentHeader());
QueryPlanStepPtr join_step
- = std::make_unique<JoinStep>(left->getCurrentDataStream(),
right->getCurrentDataStream(), hash_join, block_size, 1, false);
+ = std::make_unique<JoinStep>(left->getCurrentHeader(),
right->getCurrentHeader(), hash_join, block_size, 1, false);
std::vector<QueryPlanPtr> plans;
plans.emplace_back(std::move(left));
diff --git a/cpp-ch/local-engine/tests/gtest_ch_join.cpp
b/cpp-ch/local-engine/tests/gtest_ch_join.cpp
index af661c297f..52120cede0 100644
--- a/cpp-ch/local-engine/tests/gtest_ch_join.cpp
+++ b/cpp-ch/local-engine/tests/gtest_ch_join.cpp
@@ -109,23 +109,23 @@ TEST(TestJoin, simple)
if (right_convert_actions)
{
- auto converting_step =
std::make_unique<ExpressionStep>(right_plan.getCurrentDataStream(),
std::move(*right_convert_actions));
+ auto converting_step =
std::make_unique<ExpressionStep>(right_plan.getCurrentHeader(),
std::move(*right_convert_actions));
converting_step->setStepDescription("Convert joined columns");
right_plan.addStep(std::move(converting_step));
}
if (left_convert_actions)
{
- auto converting_step =
std::make_unique<ExpressionStep>(right_plan.getCurrentDataStream(),
std::move(*right_convert_actions));
+ auto converting_step =
std::make_unique<ExpressionStep>(right_plan.getCurrentHeader(),
std::move(*right_convert_actions));
converting_step->setStepDescription("Convert joined columns");
left_plan.addStep(std::move(converting_step));
}
- auto hash_join = std::make_shared<HashJoin>(join,
right_plan.getCurrentDataStream().header);
+ auto hash_join = std::make_shared<HashJoin>(join,
right_plan.getCurrentHeader());
QueryPlanStepPtr join_step
- = std::make_unique<JoinStep>(left_plan.getCurrentDataStream(),
right_plan.getCurrentDataStream(), hash_join, 8192, 1, false);
+ = std::make_unique<JoinStep>(left_plan.getCurrentHeader(),
right_plan.getCurrentHeader(), hash_join, 8192, 1, false);
- std::cerr << "join step:" <<
join_step->getOutputStream().header.dumpStructure() << std::endl;
+ std::cerr << "join step:" << join_step->getOutputHeader().dumpStructure()
<< std::endl;
std::vector<QueryPlanPtr> plans;
plans.emplace_back(std::make_unique<QueryPlan>(std::move(left_plan)));
@@ -133,11 +133,11 @@ TEST(TestJoin, simple)
auto query_plan = QueryPlan();
query_plan.unitePlans(std::move(join_step), {std::move(plans)});
- std::cerr << query_plan.getCurrentDataStream().header.dumpStructure() <<
std::endl;
- ActionsDAG
project{query_plan.getCurrentDataStream().header.getNamesAndTypesList()};
+ std::cerr << query_plan.getCurrentHeader().dumpStructure() << std::endl;
+ ActionsDAG project{query_plan.getCurrentHeader().getNamesAndTypesList()};
project.project(
{NameWithAlias("colA", "colA"), NameWithAlias("colB", "colB"),
NameWithAlias("colD", "colD"), NameWithAlias("colC", "colC")});
- QueryPlanStepPtr project_step =
std::make_unique<ExpressionStep>(query_plan.getCurrentDataStream(),
std::move(project));
+ QueryPlanStepPtr project_step =
std::make_unique<ExpressionStep>(query_plan.getCurrentHeader(),
std::move(project));
query_plan.addStep(std::move(project_step));
auto pipeline =
query_plan.buildQueryPipeline(QueryPlanOptimizationSettings(),
BuildQueryPipelineSettings());
auto executable_pipe =
QueryPipelineBuilder::getPipeline(std::move(*pipeline));
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]