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 3c2c51f7fb [GLUTEN-12597][CORE] Migrate ReadRel read_type to Substrait
0.98 (add iceberg_table, relocate stream_kafka) (#12832)
3c2c51f7fb is described below
commit 3c2c51f7fb6413f1bbbf739ad72811ad48202d9c
Author: Niels Pardon <[email protected]>
AuthorDate: Fri Sep 25 16:57:03 2026 +0200
[GLUTEN-12597][CORE] Migrate ReadRel read_type to Substrait 0.98 (add
iceberg_table, relocate stream_kafka) (#12832)
---
.../substrait/proto/substrait/algebra.proto | 29 +++++++++-
.../gluten/substrait/rel/ReadRelProtoSuite.scala | 64 ++++++++++++++++++++++
2 files changed, 92 insertions(+), 1 deletion(-)
diff --git
a/gluten-substrait/src/main/resources/substrait/proto/substrait/algebra.proto
b/gluten-substrait/src/main/resources/substrait/proto/substrait/algebra.proto
index ca0b15d1c3..b5bfbd703b 100644
---
a/gluten-substrait/src/main/resources/substrait/proto/substrait/algebra.proto
+++
b/gluten-substrait/src/main/resources/substrait/proto/substrait/algebra.proto
@@ -69,7 +69,11 @@ message ReadRel {
LocalFiles local_files = 6;
NamedTable named_table = 7;
ExtensionTable extension_table = 8;
- bool stream_kafka = 9;
+ IcebergTable iceberg_table = 9;
+ // Gluten addition: streaming Kafka source. Relocated from field 9 to
+ // field 1000 so the official Substrait iceberg_table can occupy its
+ // 0.98 slot.
+ bool stream_kafka = 1000;
}
// A base table. The list of string is used to represent namespacing (e.g.,
mydb.mytable).
@@ -79,6 +83,29 @@ message ReadRel {
substrait.extensions.AdvancedExtension advanced_extension = 10;
}
+ // Read an Iceberg Table
+ message IcebergTable {
+ oneof table_type {
+ MetadataFileRead direct = 1;
+ // future: add catalog table types (e.g. rest api, latest metadata in
path, etc)
+ }
+
+ // Read an Iceberg table using a metadata file. Implicit assumption:
required credentials are already known by plan consumer.
+ message MetadataFileRead {
+ // the specific uri of a metadata file (e.g.
s3://mybucket/mytable/<ver>-<uuid>.metadata.json)
+ string metadata_uri = 1;
+
+ // snapshot options. if none set, uses the current snapshot listed in
the metadata file
+ oneof snapshot {
+ // the snapshot id to read.
+ string snapshot_id = 2;
+
+ // the timestamp that should be used to select the snapshot (Time
passed in microseconds since 1970-01-01 00:00:00.000000 in UTC)
+ int64 snapshot_timestamp = 3;
+ }
+ }
+ }
+
// A table composed of expressions.
message VirtualTable {
reserved 1;
diff --git
a/gluten-substrait/src/test/scala/org/apache/gluten/substrait/rel/ReadRelProtoSuite.scala
b/gluten-substrait/src/test/scala/org/apache/gluten/substrait/rel/ReadRelProtoSuite.scala
new file mode 100644
index 0000000000..0f458692ee
--- /dev/null
+++
b/gluten-substrait/src/test/scala/org/apache/gluten/substrait/rel/ReadRelProtoSuite.scala
@@ -0,0 +1,64 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.gluten.substrait.rel
+
+import com.google.protobuf.Descriptors.Descriptor
+import io.substrait.proto.ReadRel
+import org.scalatest.funsuite.AnyFunSuite
+
+/**
+ * Pins the wire tags of the vendored `ReadRel.read_type` oneof after rebasing
it onto upstream
+ * Substrait v0.98.0: adding the official `iceberg_table = 9` and relocating
Gluten's `stream_kafka`
+ * graft off the field-9 collision into the 1000+ range. Producer and consumer
share one schema, so
+ * a renumber round-trips cleanly through the generated classes and cannot be
caught by exercising
+ * them; these assert on the descriptors instead. See
docs/developers/SubstraitModifications.md for
+ * the numbering convention.
+ */
+class ReadRelProtoSuite extends AnyFunSuite {
+
+ private def assertFieldNumbers(descriptor: Descriptor, expected: (String,
Int)*): Unit =
+ expected.foreach {
+ case (name, number) =>
+ val field = descriptor.findFieldByName(name)
+ assert(field != null, s"${descriptor.getName} has no field named
$name")
+ assert(field.getNumber === number, s"${descriptor.getName} field $name
changed its number")
+ }
+
+ test("ReadRel.read_type field numbers match upstream v0.98.0 plus the
relocated graft") {
+ assertFieldNumbers(
+ ReadRel.getDescriptor,
+ "virtual_table" -> 5,
+ "local_files" -> 6,
+ "named_table" -> 7,
+ "extension_table" -> 8,
+ // Official Substrait 0.98 addition; must own field 9.
+ "iceberg_table" -> 9,
+ // Gluten-local graft, relocated off upstream's field 9 to the 1000+
range so that
+ // iceberg_table can take its 0.98 slot.
+ "stream_kafka" -> 1000
+ )
+ }
+
+ test("ReadRel.IcebergTable structure matches the vendored upstream layout") {
+ assertFieldNumbers(ReadRel.IcebergTable.getDescriptor, "direct" -> 1)
+ assertFieldNumbers(
+ ReadRel.IcebergTable.MetadataFileRead.getDescriptor,
+ "metadata_uri" -> 1,
+ "snapshot_id" -> 2,
+ "snapshot_timestamp" -> 3)
+ }
+}
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]