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]

Reply via email to