This is an automated email from the ASF dual-hosted git repository.
zhouyuan pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/gluten.git
The following commit(s) were added to refs/heads/main by this push:
new 05ec9d1b87 [GLUTEN-12597][CORE] Make AdvancedExtension.optimization
repeated (#12642)
05ec9d1b87 is described below
commit 05ec9d1b87a1c3b09e32fd61e462e0361fdf8c6a
Author: Niels Pardon <[email protected]>
AuthorDate: Wed Jul 29 10:11:00 2026 +0200
[GLUTEN-12597][CORE] Make AdvancedExtension.optimization repeated (#12642)
Substrait 0.98 changed extensions.proto AdvancedExtension.optimization from
a singular to a repeated google.protobuf.Any (substrait-io/substrait). Update
all producers and consumers accordingly:
- Producers: AdvancedExtensionNode (JVM) and CHFormatWriterInjects use
addOptimization(...) instead of setOptimization(...); VeloxToSubstraitPlan uses
add_optimization() instead of mutable_optimization().
- Consumers: has_optimization() -> optimization_size() > 0 and
optimization() -> optimization(0) in the Velox SubstraitParser /
SubstraitToVeloxPlan and in the ClickHouse RelParsers
(Read/Cross/Aggregate/Join/GroupLimit/Write) and local_engine_jni.
The old singular optimization() accessor returned a default-constructed Any
when unset, so consumers could read optimization().value()/UnpackTo()
unconditionally. optimization(0) on an empty repeated field is out-of-bounds
and crashes. The Velox reads are all guarded by optimization_size() > 0; the
ClickHouse reads relied on that null-safe behavior, so route them through a
firstOptimizationOrDefault() helper (SubstraitParserUtils.h) that returns
Any::default_instance() when no optimiz [...]
Validated locally: proto (protoc), gluten-substrait JVM build, and Velox
native build (libgluten/libvelox). The ClickHouse native parsers are covered by
CI (local libch build requires a Linux/Docker toolchain). Part of #12597.
Depends on #12604 (URI->URN).
---
.../sql/execution/datasources/v1/CHFormatWriterInjects.scala | 2 +-
cpp-ch/local-engine/Parser/RelParsers/AggregateRelParser.cpp | 2 +-
cpp-ch/local-engine/Parser/RelParsers/CrossRelParser.cpp | 3 ++-
cpp-ch/local-engine/Parser/RelParsers/GroupLimitRelParser.cpp | 8 ++++----
cpp-ch/local-engine/Parser/RelParsers/JoinRelParser.cpp | 2 +-
cpp-ch/local-engine/Parser/RelParsers/ReadRelParser.cpp | 4 ++--
cpp-ch/local-engine/Parser/RelParsers/WriteRelParser.cpp | 3 ++-
cpp-ch/local-engine/Parser/SubstraitParserUtils.h | 11 +++++++++++
cpp-ch/local-engine/local_engine_jni.cpp | 4 ++--
cpp/velox/substrait/SubstraitParser.cc | 8 ++++----
cpp/velox/substrait/SubstraitToVeloxPlan.cc | 4 ++--
cpp/velox/substrait/VeloxToSubstraitPlan.cc | 2 +-
.../gluten/substrait/extensions/AdvancedExtensionNode.java | 2 +-
.../substrait/proto/substrait/extensions/extensions.proto | 2 +-
14 files changed, 35 insertions(+), 22 deletions(-)
diff --git
a/backends-clickhouse/src/main/scala/org/apache/spark/sql/execution/datasources/v1/CHFormatWriterInjects.scala
b/backends-clickhouse/src/main/scala/org/apache/spark/sql/execution/datasources/v1/CHFormatWriterInjects.scala
index 68c7cbea9b..0959b038c0 100644
---
a/backends-clickhouse/src/main/scala/org/apache/spark/sql/execution/datasources/v1/CHFormatWriterInjects.scala
+++
b/backends-clickhouse/src/main/scala/org/apache/spark/sql/execution/datasources/v1/CHFormatWriterInjects.scala
@@ -62,7 +62,7 @@ trait CHFormatWriterInjects extends
GlutenFormatWriterInjectsBase {
.setAdvancedExtension(
AdvancedExtension
.newBuilder()
- .setOptimization(Any.pack(createNativeWrite(outputPath,
context)))
+ .addOptimization(Any.pack(createNativeWrite(outputPath,
context)))
.build())
.build())
.build()
diff --git a/cpp-ch/local-engine/Parser/RelParsers/AggregateRelParser.cpp
b/cpp-ch/local-engine/Parser/RelParsers/AggregateRelParser.cpp
index f81d8ec77b..6880dc8d2c 100644
--- a/cpp-ch/local-engine/Parser/RelParsers/AggregateRelParser.cpp
+++ b/cpp-ch/local-engine/Parser/RelParsers/AggregateRelParser.cpp
@@ -129,7 +129,7 @@ void AggregateRelParser::setup(DB::QueryPlanPtr query_plan,
const substrait::Rel
has_complete_stage =
phase_set.contains(substrait::AggregationPhase::AGGREGATION_PHASE_INITIAL_TO_RESULT);
bool next_step_is_agg = false;
google::protobuf::StringValue raw_extra_params;
-
raw_extra_params.ParseFromString(aggregate_rel->advanced_extension().optimization().value());
+
raw_extra_params.ParseFromString(firstOptimizationOrDefault(aggregate_rel->advanced_extension()).value());
auto extra_params =
AggregateOptimizationInfo::parse(raw_extra_params.value());
if (aggregate_rel->measures().empty())
{
diff --git a/cpp-ch/local-engine/Parser/RelParsers/CrossRelParser.cpp
b/cpp-ch/local-engine/Parser/RelParsers/CrossRelParser.cpp
index 6442abb3ed..854752c555 100644
--- a/cpp-ch/local-engine/Parser/RelParsers/CrossRelParser.cpp
+++ b/cpp-ch/local-engine/Parser/RelParsers/CrossRelParser.cpp
@@ -27,6 +27,7 @@
#include <Parser/AdvancedParametersParseUtil.h>
#include <Parser/ExpressionParser.h>
#include <Parser/SerializedPlanParser.h>
+#include <Parser/SubstraitParserUtils.h>
#include <Parsers/ASTIdentifier.h>
#include <Processors/QueryPlan/ExpressionStep.h>
#include <Processors/QueryPlan/FilterStep.h>
@@ -168,7 +169,7 @@ void CrossRelParser::renamePlanColumns(DB::QueryPlan &
left, DB::QueryPlan & rig
DB::QueryPlanPtr CrossRelParser::parseJoin(const substrait::CrossRel & join,
DB::QueryPlanPtr left, DB::QueryPlanPtr right)
{
google::protobuf::StringValue optimization_info;
-
optimization_info.ParseFromString(join.advanced_extension().optimization().value());
+
optimization_info.ParseFromString(firstOptimizationOrDefault(join.advanced_extension()).value());
auto join_opt_info =
JoinOptimizationInfo::parse(optimization_info.value());
const auto & storage_join_key = join_opt_info.storage_join_key;
auto storage_join = !storage_join_key.empty() ?
BroadcastJoinBuilder::getJoin(storage_join_key) : nullptr;
diff --git a/cpp-ch/local-engine/Parser/RelParsers/GroupLimitRelParser.cpp
b/cpp-ch/local-engine/Parser/RelParsers/GroupLimitRelParser.cpp
index 2d49a8b453..e76d1ec7a3 100644
--- a/cpp-ch/local-engine/Parser/RelParsers/GroupLimitRelParser.cpp
+++ b/cpp-ch/local-engine/Parser/RelParsers/GroupLimitRelParser.cpp
@@ -80,7 +80,7 @@ GroupLimitRelParser::parse(DB::QueryPlanPtr current_plan_,
const substrait::Rel
{
const auto win_rel_def = rel.windowgrouplimit();
google::protobuf::StringValue optimize_info_str;
-
optimize_info_str.ParseFromString(win_rel_def.advanced_extension().optimization().value());
+
optimize_info_str.ParseFromString(firstOptimizationOrDefault(win_rel_def.advanced_extension()).value());
auto optimization_info =
WindowGroupOptimizationInfo::parse(optimize_info_str.value());
if (optimization_info.is_aggregate_group_limit)
{
@@ -134,7 +134,7 @@ WindowGroupLimitRelParser::parse(DB::QueryPlanPtr
current_plan_, const substrait
{
const auto win_rel_def = rel.windowgrouplimit();
google::protobuf::StringValue optimize_info_str;
-
optimize_info_str.ParseFromString(win_rel_def.advanced_extension().optimization().value());
+
optimize_info_str.ParseFromString(firstOptimizationOrDefault(win_rel_def.advanced_extension()).value());
auto optimization_info =
WindowGroupOptimizationInfo::parse(optimize_info_str.value());
window_function_name = optimization_info.window_function;
@@ -198,7 +198,7 @@ DB::QueryPlanPtr AggregateGroupLimitRelParser::parse(
win_rel_def = &rel.windowgrouplimit();
google::protobuf::StringValue optimize_info_str;
-
optimize_info_str.ParseFromString(win_rel_def->advanced_extension().optimization().value());
+
optimize_info_str.ParseFromString(firstOptimizationOrDefault(win_rel_def->advanced_extension()).value());
auto optimization_info =
WindowGroupOptimizationInfo::parse(optimize_info_str.value());
limit = static_cast<size_t>(win_rel_def->limit());
aggregate_function_name =
getAggregateFunctionName(optimization_info.window_function);
@@ -499,7 +499,7 @@ static DB::WindowFunctionDescription
buildWindowFunctionDescription(const std::s
void AggregateGroupLimitRelParser::addWindowLimitStep(DB::QueryPlan & plan)
{
google::protobuf::StringValue optimize_info_str;
-
optimize_info_str.ParseFromString(win_rel_def->advanced_extension().optimization().value());
+
optimize_info_str.ParseFromString(firstOptimizationOrDefault(win_rel_def->advanced_extension()).value());
auto optimization_info =
WindowGroupOptimizationInfo::parse(optimize_info_str.value());
auto window_function_name = optimization_info.window_function;
diff --git a/cpp-ch/local-engine/Parser/RelParsers/JoinRelParser.cpp
b/cpp-ch/local-engine/Parser/RelParsers/JoinRelParser.cpp
index 605f3aea5f..05ae167249 100644
--- a/cpp-ch/local-engine/Parser/RelParsers/JoinRelParser.cpp
+++ b/cpp-ch/local-engine/Parser/RelParsers/JoinRelParser.cpp
@@ -204,7 +204,7 @@ DB::QueryPlanPtr JoinRelParser::parseJoin(const
substrait::JoinRel & join, DB::Q
{
auto join_config = JoinConfig::loadFromContext(getContext());
google::protobuf::StringValue optimization_info;
-
optimization_info.ParseFromString(join.advanced_extension().optimization().value());
+
optimization_info.ParseFromString(firstOptimizationOrDefault(join.advanced_extension()).value());
auto join_opt_info =
JoinOptimizationInfo::parse(optimization_info.value());
LOG_DEBUG(getLogger("JoinRelParser"), "optimization info:{}",
optimization_info.value());
auto storage_join = join_opt_info.is_broadcast ?
BroadcastJoinBuilder::getJoin(join_opt_info.storage_join_key) : nullptr;
diff --git a/cpp-ch/local-engine/Parser/RelParsers/ReadRelParser.cpp
b/cpp-ch/local-engine/Parser/RelParsers/ReadRelParser.cpp
index 9819dc2c01..17bdbc5cfd 100644
--- a/cpp-ch/local-engine/Parser/RelParsers/ReadRelParser.cpp
+++ b/cpp-ch/local-engine/Parser/RelParsers/ReadRelParser.cpp
@@ -131,7 +131,7 @@ bool ReadRelParser::isReadRelFromMergeTree(const
substrait::ReadRel & rel)
return false;
google::protobuf::StringValue optimization;
-
optimization.ParseFromString(rel.advanced_extension().optimization().value());
+
optimization.ParseFromString(firstOptimizationOrDefault(rel.advanced_extension()).value());
ReadBufferFromString in(optimization.value());
if (!checkString("isMergeTree=", in))
return false;
@@ -148,7 +148,7 @@ bool ReadRelParser::isReadRelFromRange(const
substrait::ReadRel & rel)
return false;
google::protobuf::StringValue optimization;
-
optimization.ParseFromString(rel.advanced_extension().optimization().value());
+
optimization.ParseFromString(firstOptimizationOrDefault(rel.advanced_extension()).value());
ReadBufferFromString in(optimization.value());
if (!checkString("isRange=", in))
return false;
diff --git a/cpp-ch/local-engine/Parser/RelParsers/WriteRelParser.cpp
b/cpp-ch/local-engine/Parser/RelParsers/WriteRelParser.cpp
index aaa0779fe7..48d835d6b2 100644
--- a/cpp-ch/local-engine/Parser/RelParsers/WriteRelParser.cpp
+++ b/cpp-ch/local-engine/Parser/RelParsers/WriteRelParser.cpp
@@ -21,6 +21,7 @@
#include <DataTypes/DataTypeTuple.h>
#include <Interpreters/Context.h>
#include <Interpreters/ExpressionActions.h>
+#include <Parser/SubstraitParserUtils.h>
#include <Parser/TypeParser.h>
#include <Processors/Transforms/ExpressionTransform.h>
#include <Processors/Transforms/MaterializingTransform.h>
@@ -196,7 +197,7 @@ void addSinkTransform(const DB::ContextPtr & context, const
substrait::WriteRel
const substrait::NamedObjectWrite & named_table = write_rel.named_table();
local_engine::Write write;
- if (!named_table.advanced_extension().optimization().UnpackTo(&write))
+ if
(!firstOptimizationOrDefault(named_table.advanced_extension()).UnpackTo(&write))
throw DB::Exception(DB::ErrorCodes::LOGICAL_ERROR, "Failed to unpack
write optimization with local_engine::Write.");
assert(write.has_common());
const substrait::NamedStruct & table_schema = write_rel.table_schema();
diff --git a/cpp-ch/local-engine/Parser/SubstraitParserUtils.h
b/cpp-ch/local-engine/Parser/SubstraitParserUtils.h
index a9e21ba579..836b0a047c 100644
--- a/cpp-ch/local-engine/Parser/SubstraitParserUtils.h
+++ b/cpp-ch/local-engine/Parser/SubstraitParserUtils.h
@@ -78,6 +78,17 @@ inline std::string toString(const google::protobuf::Any &
any)
return sv.value();
}
+/// Substrait 0.98 changed AdvancedExtension.optimization from a singular to a
repeated field.
+/// The old singular accessor returned a default-constructed Any when the
field was unset (empty
+/// value(), empty type_url), so callers could unconditionally read
optimization().value()/UnpackTo().
+/// optimization(0) on an empty repeated field is out-of-bounds and crashes,
so route reads of the
+/// first optimization through this helper to preserve the null-safe behavior.
+inline const google::protobuf::Any & firstOptimizationOrDefault(const
substrait::extensions::AdvancedExtension & advanced_extension)
+{
+ return advanced_extension.optimization_size() > 0 ?
advanced_extension.optimization(0)
+ :
google::protobuf::Any::default_instance();
+}
+
namespace SubstraitParserUtils
{
std::optional<size_t> getStructFieldIndex(const substrait::Expression & e);
diff --git a/cpp-ch/local-engine/local_engine_jni.cpp
b/cpp-ch/local-engine/local_engine_jni.cpp
index 654078717f..12b9d2a1f4 100644
--- a/cpp-ch/local-engine/local_engine_jni.cpp
+++ b/cpp-ch/local-engine/local_engine_jni.cpp
@@ -961,7 +961,7 @@ JNIEXPORT jlong
Java_org_apache_spark_sql_execution_datasources_CHDatasourceJniW
assert(write_rel.has_named_table());
const substrait::NamedObjectWrite & named_table = write_rel.named_table();
local_engine::Write write_opt;
- named_table.advanced_extension().optimization().UnpackTo(&write_opt);
+
local_engine::firstOptimizationOrDefault(named_table.advanced_extension()).UnpackTo(&write_opt);
DB::Block preferred_schema =
local_engine::TypeParser::buildBlockFromNamedStructWithoutDFS(write_rel.table_schema());
const auto file_uri = jstring2string(env, file_uri_);
@@ -990,7 +990,7 @@ JNIEXPORT jlong
Java_org_apache_spark_sql_execution_datasources_CHDatasourceJniW
assert(write_rel.has_named_table());
const substrait::NamedObjectWrite & named_table = write_rel.named_table();
local_engine::Write write;
- if (!named_table.advanced_extension().optimization().UnpackTo(&write))
+ if
(!local_engine::firstOptimizationOrDefault(named_table.advanced_extension()).UnpackTo(&write))
throw DB::Exception(DB::ErrorCodes::LOGICAL_ERROR, "Failed to unpack
write optimization with local_engine::Write.");
assert(write.has_common());
assert(write.has_mergetree());
diff --git a/cpp/velox/substrait/SubstraitParser.cc
b/cpp/velox/substrait/SubstraitParser.cc
index a2e76a3131..54bf8d4f24 100644
--- a/cpp/velox/substrait/SubstraitParser.cc
+++ b/cpp/velox/substrait/SubstraitParser.cc
@@ -279,9 +279,9 @@ std::string SubstraitParser::mapToVeloxFunction(const
std::string& substraitFunc
bool SubstraitParser::configSetInOptimization(
const ::substrait::extensions::AdvancedExtension& extension,
const std::string& config) {
- if (extension.has_optimization()) {
+ if (extension.optimization_size() > 0) {
google::protobuf::StringValue msg;
- extension.optimization().UnpackTo(&msg);
+ extension.optimization(0).UnpackTo(&msg);
std::size_t pos = msg.value().find(config);
if ((pos != std::string::npos) && (msg.value().substr(pos + config.size(),
1) == "1")) {
return true;
@@ -294,9 +294,9 @@ bool SubstraitParser::checkWindowFunction(
const ::substrait::extensions::AdvancedExtension& extension,
const std::string& targetFunction) {
const std::string config = "window_function=";
- if (extension.has_optimization()) {
+ if (extension.optimization_size() > 0) {
google::protobuf::StringValue msg;
- extension.optimization().UnpackTo(&msg);
+ extension.optimization(0).UnpackTo(&msg);
std::size_t pos = msg.value().find(config);
if ((pos != std::string::npos) && (msg.value().size() >=
targetFunction.size()) &&
(msg.value().substr(pos + config.size(), targetFunction.size()) ==
targetFunction)) {
diff --git a/cpp/velox/substrait/SubstraitToVeloxPlan.cc
b/cpp/velox/substrait/SubstraitToVeloxPlan.cc
index fedd950374..9e85daa206 100644
--- a/cpp/velox/substrait/SubstraitToVeloxPlan.cc
+++ b/cpp/velox/substrait/SubstraitToVeloxPlan.cc
@@ -870,8 +870,8 @@ core::PlanNodePtr
SubstraitToVeloxPlanConverter::toVeloxPlan(const ::substrait::
GLUTEN_CHECK(writeRel.named_table().has_advanced_extension(), "Advanced
extension not found in WriteRel");
const auto& ext = writeRel.named_table().advanced_extension();
- GLUTEN_CHECK(ext.has_optimization(), "Extension optimization not found in
WriteRel");
- const auto& opt = ext.optimization();
+ GLUTEN_CHECK(ext.optimization_size() > 0, "Extension optimization not found
in WriteRel");
+ const auto& opt = ext.optimization(0);
gluten::ConfigMap confMap;
opt.UnpackTo(&confMap);
std::unordered_map<std::string, std::string> writeConfs;
diff --git a/cpp/velox/substrait/VeloxToSubstraitPlan.cc
b/cpp/velox/substrait/VeloxToSubstraitPlan.cc
index 39d9d2e152..a137872411 100644
--- a/cpp/velox/substrait/VeloxToSubstraitPlan.cc
+++ b/cpp/velox/substrait/VeloxToSubstraitPlan.cc
@@ -282,7 +282,7 @@ void VeloxToSubstraitPlanConvertor::toSubstrait(
substrait::extensions::AdvancedExtension ae{};
google::protobuf::StringValue msg;
msg.set_value("allowFlush=1");
- ae.mutable_optimization()->PackFrom(msg);
+ ae.add_optimization()->PackFrom(msg);
aggregateRel->mutable_advanced_extension()->MergeFrom(ae);
break;
}
diff --git
a/gluten-substrait/src/main/java/org/apache/gluten/substrait/extensions/AdvancedExtensionNode.java
b/gluten-substrait/src/main/java/org/apache/gluten/substrait/extensions/AdvancedExtensionNode.java
index b75f4e587f..df9288cf79 100644
---
a/gluten-substrait/src/main/java/org/apache/gluten/substrait/extensions/AdvancedExtensionNode.java
+++
b/gluten-substrait/src/main/java/org/apache/gluten/substrait/extensions/AdvancedExtensionNode.java
@@ -51,7 +51,7 @@ public class AdvancedExtensionNode implements Serializable {
public AdvancedExtension toProtobuf() {
AdvancedExtension.Builder extensionBuilder =
AdvancedExtension.newBuilder();
if (optimization != null) {
- extensionBuilder.setOptimization(optimization);
+ extensionBuilder.addOptimization(optimization);
}
if (enhancement != null) {
extensionBuilder.setEnhancement(enhancement);
diff --git
a/gluten-substrait/src/main/resources/substrait/proto/substrait/extensions/extensions.proto
b/gluten-substrait/src/main/resources/substrait/proto/substrait/extensions/extensions.proto
index 1e24ace044..105b58bdb9 100644
---
a/gluten-substrait/src/main/resources/substrait/proto/substrait/extensions/extensions.proto
+++
b/gluten-substrait/src/main/resources/substrait/proto/substrait/extensions/extensions.proto
@@ -82,7 +82,7 @@ message SimpleExtensionDeclaration {
message AdvancedExtension {
// An optimization is helpful information that don't influence semantics. May
// be ignored by a consumer.
- google.protobuf.Any optimization = 1;
+ repeated google.protobuf.Any optimization = 1;
// An enhancement alter semantics. Cannot be ignored by a consumer.
google.protobuf.Any enhancement = 2;
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]