Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
@@ -0,0 +1,84 @@
// 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.doris.iceberg;

import org.apache.iceberg.Schema;

import java.io.ByteArrayInputStream;
import java.io.IOException;
import java.io.ObjectInputStream;
import java.io.ObjectStreamClass;
import java.io.UncheckedIOException;
import java.nio.charset.StandardCharsets;
import java.util.Base64;

final class IcebergSerializationCompat {
private static final long ICEBERG_1_10_1_SCHEMA_UID = 6812231194765760118L;
private static final long ICEBERG_1_11_0_SCHEMA_UID = -1265875184407129845L;
private static final ObjectStreamClass LOCAL_SCHEMA_DESCRIPTOR = ObjectStreamClass.lookup(Schema.class);

private IcebergSerializationCompat() {
}

@SuppressWarnings({"DangerousJavaDeserialization", "unchecked"})
static <T> T deserializeFromBase64(String base64) {
if (base64 == null) {
return null;
}
byte[] bytes = Base64.getMimeDecoder().decode(base64.getBytes(StandardCharsets.UTF_8));
try (ByteArrayInputStream input = new ByteArrayInputStream(bytes);
ObjectInputStream objectInput = new SchemaCompatibleObjectInputStream(input)) {
return (T) objectInput.readObject();
} catch (IOException e) {
throw new UncheckedIOException("Failed to deserialize object", e);
} catch (ClassNotFoundException e) {
throw new RuntimeException("Could not read object ", e);
}
}

private static final class SchemaCompatibleObjectInputStream extends ObjectInputStream {
private SchemaCompatibleObjectInputStream(ByteArrayInputStream input) throws IOException {
super(input);
}

@Override
protected ObjectStreamClass readClassDescriptor() throws IOException, ClassNotFoundException {
ObjectStreamClass descriptor = super.readClassDescriptor();
// Iceberg 1.11 changed Schema's generated UID without changing its serialized fields. Accept only
// that known rolling-upgrade pair so unrelated or future class-layout changes still fail closed.
if (Schema.class.getName().equals(descriptor.getName())
&& descriptor.getSerialVersionUID() == ICEBERG_1_10_1_SCHEMA_UID
&& LOCAL_SCHEMA_DESCRIPTOR.getSerialVersionUID() == ICEBERG_1_11_0_SCHEMA_UID) {
return LOCAL_SCHEMA_DESCRIPTOR;
}
return descriptor;
}

@Override
protected Class<?> resolveClass(ObjectStreamClass descriptor) throws IOException, ClassNotFoundException {
String className = descriptor.getName();
if (className.indexOf('/') >= 0) {
// Replacing Schema's stream descriptor exposes JVM-style array signatures with slashes; resolve
// them through the scanner's isolated class loader after normalizing to Class.forName syntax.
return Class.forName(className.replace('/', '.'), false,
IcebergSerializationCompat.class.getClassLoader());
}
return super.resolveClass(descriptor);
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -28,7 +28,6 @@
import org.apache.iceberg.FileScanTask;
import org.apache.iceberg.StructLike;
import org.apache.iceberg.io.CloseableIterator;
import org.apache.iceberg.util.SerializationUtil;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

Expand All @@ -55,7 +54,7 @@ public IcebergSysTableJniScanner(int batchSize, Map<String, String> params) {
String serializedSplitParams = params.get("serialized_split");
Preconditions.checkArgument(serializedSplitParams != null && !serializedSplitParams.isEmpty(),
"serialized_split should not be empty");
this.scanTask = SerializationUtil.deserializeFromBase64(serializedSplitParams);
this.scanTask = IcebergSerializationCompat.deserializeFromBase64(serializedSplitParams);
String requiredFieldsParam = params.get("required_fields");
Preconditions.checkArgument(requiredFieldsParam != null && !requiredFieldsParam.isEmpty(),
"required_fields should not be empty");
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,76 @@
// 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.doris.iceberg;

import org.apache.iceberg.FileScanTask;
import org.apache.iceberg.Schema;
import org.apache.iceberg.types.Types;
import org.junit.Assert;
import org.junit.Test;

import java.io.ByteArrayOutputStream;
import java.io.IOException;
import java.io.ObjectOutputStream;
import java.util.Base64;

public class IcebergSerializationCompatTest {
// An empty StaticDataTask serialized by Iceberg 1.10.1. It carries the old Schema serialVersionUID while
// keeping all data neutral and local, so the fixture remains stable and safe to commit.
private static final String ICEBERG_1_10_1_TASK = "rO0ABXNyACFvcmcuYXBhY2hlLmljZWJlcmcuU3RhdGljRGF0YVRhc2t2PvVIpr/rlAIABEwADG1ldGFkYXRh"
+ "RmlsZXQAHUxvcmcvYXBhY2hlL2ljZWJlcmcvRGF0YUZpbGU7TAAPcHJvamVjdGVkU2NoZW1hdAAbTG9yZy9h"
+ "cGFjaGUvaWNlYmVyZy9TY2hlbWE7WwAEcm93c3QAIFtMb3JnL2FwYWNoZS9pY2ViZXJnL1N0cnVjdExpa2U7"
+ "TAALdGFibGVTY2hlbWFxAH4AAnhwcHNyABlvcmcuYXBhY2hlLmljZWJlcmcuU2NoZW1hXonoLcvZFnYCAARJ"
+ "AA5oaWdoZXN0RmllbGRJZEkACHNjaGVtYUlkWwASaWRlbnRpZmllckZpZWxkSWRzdAACW0lMAAZzdHJ1Y3R0"
+ "ACtMb3JnL2FwYWNoZS9pY2ViZXJnL3R5cGVzL1R5cGVzJFN0cnVjdFR5cGU7eHAAAAABAAAAAHVyAAJbSU26"
+ "YCZ26rKlAgAAeHAAAAAAc3IAKW9yZy5hcGFjaGUuaWNlYmVyZy50eXBlcy5UeXBlcyRTdHJ1Y3RUeXBlY2OW"
+ "YF+O53QCAAFbAAZmaWVsZHN0AC1bTG9yZy9hcGFjaGUvaWNlYmVyZy90eXBlcy9UeXBlcyROZXN0ZWRGaWVs"
+ "ZDt4cgAob3JnLmFwYWNoZS5pY2ViZXJnLnR5cGVzLlR5cGUkTmVzdGVkVHlwZWpUXj112XcBAgAAeHB1cgAt"
+ "W0xvcmcuYXBhY2hlLmljZWJlcmcudHlwZXMuVHlwZXMkTmVzdGVkRmllbGQ7A428r/O1h1gCAAB4cAAAAAFz"
+ "cgAqb3JnLmFwYWNoZS5pY2ViZXJnLnR5cGVzLlR5cGVzJE5lc3RlZEZpZWxkRnIcmOgj/wICAAdJAAJpZFoA"
+ "CmlzT3B0aW9uYWxMAANkb2N0ABJMamF2YS9sYW5nL1N0cmluZztMAA5pbml0aWFsRGVmYXVsdHQAKExvcmcv"
+ "YXBhY2hlL2ljZWJlcmcvZXhwcmVzc2lvbnMvTGl0ZXJhbDtMAARuYW1lcQB+ABJMAAR0eXBldAAfTG9yZy9h"
+ "cGFjaGUvaWNlYmVyZy90eXBlcy9UeXBlO0wADHdyaXRlRGVmYXVsdHEAfgATeHAAAAABAHBwdAACaWRzcgAs"
+ "b3JnLmFwYWNoZS5pY2ViZXJnLnR5cGVzLlByaW1pdGl2ZUxpa2VIb2xkZXLbKTPuyM89cAIAAUwADHR5cGVB"
+ "c1N0cmluZ3EAfgASeHB0AANpbnRwdXIAIFtMb3JnL2FwYWNoZS5pY2ViZXJnLlN0cnVjdExpa2U7MJKGbVap"
+ "uFMCAAB4cAAAAABxAH4ACA==";

@Test
public void deserializesIceberg1101SystemTableTask() {
FileScanTask task = IcebergSerializationCompat.deserializeFromBase64(ICEBERG_1_10_1_TASK);

Assert.assertEquals("id", task.schema().columns().get(0).name());
}

@Test
public void preservesCurrentIcebergSerialization() throws IOException {
Schema expected = new Schema(Types.NestedField.required(1, "id", Types.IntegerType.get()));

Schema actual = IcebergSerializationCompat.deserializeFromBase64(
serializeToBase64(expected));

Assert.assertTrue(expected.sameSchema(actual));
}

private static String serializeToBase64(Object value) throws IOException {
ByteArrayOutputStream output = new ByteArrayOutputStream();
try (ObjectOutputStream objectOutput = new ObjectOutputStream(output)) {
objectOutput.writeObject(value);
}
return Base64.getEncoder().encodeToString(output.toByteArray());
}
}
2 changes: 2 additions & 0 deletions fe/check/checkstyle/suppressions.xml
Original file line number Diff line number Diff line change
Expand Up @@ -81,6 +81,8 @@ under the License.

<!-- ignore iceberg delete file index copied from iceberg/DeleteFileIndex.java -->
<suppress files="org[\\/]apache[\\/]iceberg[\\/]DeleteFileIndex\.java" checks="[a-zA-Z0-9]*"/>
<!-- Iceberg package access is required to preserve historical schema binding in DataTableScan. -->
<suppress files="org[\\/]apache[\\/]iceberg[\\/]SchemaAwareDataTableScan\.java" checks="[a-zA-Z0-9]*"/>

<!-- ignore gensrc/thrift/ExternalTableSchema.thrift -->
<suppress files=".*thrift/schema/external/.*" checks=".*"/>
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -20,16 +20,19 @@
import java.util.Objects;

/**
* Immutable cache key for {@link ConnectorMetadataCache}: {@code (db, table, snapshotId, schemaId)}.
* Immutable cache key for {@link ConnectorMetadataCache}:
* {@code (db, table, snapshotId, schemaId, metadataGeneration)}.
*
* <p>Engine-agnostic (external-partition-derived-cache design doc §5, "cache A"): a table's derived partition
* view is a pure function of its identity plus the MVCC coordinate it was read at, so pinning that coordinate
* into the key is what makes the cache "always correct" — a new snapshot/schema yields a new key, never a stale
* hit. Non-MVCC engines (hive) or engines without a separate schema version pass {@code snapshotId = -1} /
* {@code schemaId = -1}; the key still holds them, it just means "unversioned" for that axis.
* {@code schemaId = -1}; the key still holds them, it just means "unversioned" for that axis. Connectors whose
* derived view depends on another independently evolving metadata generation may use
* {@code metadataGeneration}; other connectors leave it at {@code -1} through the four-argument constructor.
*
* <p>{@link #matches} / {@link #matchesDb} back {@link ConnectorMetadataCache#invalidateTable} /
* {@link ConnectorMetadataCache#invalidateDb}, which must drop every snapshot/schema of a (db, table) or
* {@link ConnectorMetadataCache#invalidateDb}, which must drop every metadata generation of a (db, table) or
* every table of a db — mirrors the {@code matches}/{@code matchesDb} helpers on the sibling connector caches
* ({@code MaxComputePartitionCache.PartitionKey}, {@code HiveFileListingCache.FileListingKey}).
*/
Expand All @@ -38,12 +41,18 @@ public final class ConnectorTableKey {
private final String table;
private final long snapshotId;
private final long schemaId;
private final long metadataGeneration;

public ConnectorTableKey(String db, String table, long snapshotId, long schemaId) {
this(db, table, snapshotId, schemaId, -1L);
}

public ConnectorTableKey(String db, String table, long snapshotId, long schemaId, long metadataGeneration) {
this.db = db;
this.table = table;
this.snapshotId = snapshotId;
this.schemaId = schemaId;
this.metadataGeneration = metadataGeneration;
}

public String getDb() {
Expand All @@ -62,12 +71,16 @@ public long getSchemaId() {
return schemaId;
}

/** Whether this key belongs to the given (db, table), regardless of snapshotId/schemaId. */
public long getMetadataGeneration() {
return metadataGeneration;
}

/** Whether this key belongs to the given (db, table), regardless of its metadata coordinates. */
public boolean matches(String db, String table) {
return Objects.equals(this.db, db) && Objects.equals(this.table, table);
}

/** Whether this key belongs to the given db, regardless of table/snapshotId/schemaId. */
/** Whether this key belongs to the given db, regardless of table or metadata coordinates. */
public boolean matchesDb(String db) {
return Objects.equals(this.db, db);
}
Expand All @@ -83,18 +96,20 @@ public boolean equals(Object o) {
ConnectorTableKey that = (ConnectorTableKey) o;
return snapshotId == that.snapshotId
&& schemaId == that.schemaId
&& metadataGeneration == that.metadataGeneration
&& Objects.equals(db, that.db)
&& Objects.equals(table, that.table);
}

@Override
public int hashCode() {
return Objects.hash(db, table, snapshotId, schemaId);
return Objects.hash(db, table, snapshotId, schemaId, metadataGeneration);
}

@Override
public String toString() {
return "ConnectorTableKey{db=" + db + ", table=" + table
+ ", snapshotId=" + snapshotId + ", schemaId=" + schemaId + '}';
+ ", snapshotId=" + snapshotId + ", schemaId=" + schemaId
+ ", metadataGeneration=" + metadataGeneration + '}';
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -41,6 +41,11 @@ private static ConnectorTableKey key(String db, String table, long snapshotId, l
return new ConnectorTableKey(db, table, snapshotId, schemaId);
}

private static ConnectorTableKey key(
String db, String table, long snapshotId, long schemaId, long metadataGeneration) {
return new ConnectorTableKey(db, table, snapshotId, schemaId, metadataGeneration);
}

private static ConnectorMetadataCache<String> newCache() {
return new ConnectorMetadataCache<>(ENGINE, "partition_view", new HashMap<>());
}
Expand Down Expand Up @@ -108,6 +113,23 @@ public void differentSchemaIdIsADistinctEntry() {
Assertions.assertEquals(2, loads.get(), "distinct schemaId must trigger a distinct load");
}

@Test
public void differentMetadataGenerationIsADistinctEntry() {
AtomicInteger loads = new AtomicInteger();
ConnectorMetadataCache<String> cache = newCache();

cache.get(key("db", "t", 1L, 1L, 10L), () -> "generation-10");
String second = cache.get(key("db", "t", 1L, 1L, 11L), () -> {
loads.incrementAndGet();
return "generation-11";
});

// Some table metadata evolves independently of both snapshot and schema; treating this axis as part of
// key identity prevents a connector from serving a derived view built for the preceding generation.
Assertions.assertEquals("generation-11", second);
Assertions.assertEquals(1, loads.get(), "a distinct metadata generation must trigger a distinct load");
}

@Test
public void invalidateTableEvictsAllSnapshotsOfThatTableOnly() {
AtomicInteger loads = new AtomicInteger();
Expand Down
7 changes: 4 additions & 3 deletions fe/fe-connector/fe-connector-hms-hive-shade/pom.xml
Original file line number Diff line number Diff line change
Expand Up @@ -211,17 +211,18 @@ under the License.
</exclusions>
</dependency>

<!-- iceberg-hive-metastore 1.10.1 supplies org.apache.iceberg.hive.HiveCatalog (the hms flavor);
<!-- iceberg-hive-metastore supplies org.apache.iceberg.hive.HiveCatalog (the hms flavor);
only fe-connector-iceberg loads it. iceberg-core/api/common/bundled-guava + caffeine + slf4j
are excluded — the iceberg plugin already ships iceberg-core 1.10.1 (which brings api/common/
are excluded — the iceberg plugin already ships the matching iceberg-core version (which brings api/common/
guava/caffeine) as its own child-first jar, so re-bundling them here would create duplicate
classes. Only the org.apache.iceberg.hive.* classes are kept, and the relocation below rewrites
their org.apache.thrift refs to the shared private prefix so they link against the same
relocated HiveMetaStoreClient + libthrift as the plain-HMS path. -->
<dependency>
<groupId>org.apache.iceberg</groupId>
<artifactId>iceberg-hive-metastore</artifactId>
<version>1.10.1</version>
<!-- Keep HiveCatalog bytecode aligned with the Iceberg core loaded by the connector. -->
<version>${iceberg.version}</version>
<optional>true</optional>
<exclusions>
<exclusion><groupId>org.apache.iceberg</groupId><artifactId>iceberg-core</artifactId></exclusion>
Expand Down
Loading
Loading