Copilot commented on code in PR #3601:
URL: https://github.com/apache/fluss/pull/3601#discussion_r3607496319
##########
fluss-lake/fluss-lake-iceberg/src/test/java/org/apache/fluss/lake/iceberg/IcebergLakeCatalogTest.java:
##########
@@ -1131,6 +1139,469 @@ void testAlterTableComplexTypeMismatch() {
.hasMessageContaining("Iceberg schema is not compatible with
Fluss schema");
}
+ /**
+ * If the Iceberg table already exists with a compatible schema,
re-creating the same Fluss
+ * table should succeed idempotently as long as the caller does not claim
it is the first
+ * creation of the Fluss table.
+ */
+ @Test
+ void testCreateTableIdempotentWhenCompatibleLakeTableExists() {
+ String database = "reuse_existing_db";
+ String tableName = "reuse_existing_table";
+ TablePath tablePath = TablePath.of(database, tableName);
+
+ TableDescriptor td = getTableDescriptor(FLUSS_SCHEMA);
+
+ // First creation should succeed
+ flussIcebergCatalog.createTable(tablePath, td, new
TestingLakeCatalogContext());
+
+ // Second creation with the same schema and isCreatingFlussTable=false
should succeed
+ // idempotently.
+ flussIcebergCatalog.createTable(tablePath, td, new
TestingLakeCatalogContext());
+
+ Table table =
+ flussIcebergCatalog
+ .getIcebergCatalog()
+ .loadTable(TableIdentifier.of(database, tableName));
+ assertThat(table).isNotNull();
+ }
+
+ /**
+ * When creating a fresh Fluss table, if a compatible but empty Iceberg
table already exists,
+ * the creation should still succeed (allowing a Fluss table to bind to a
pre-created empty lake
+ * table).
+ */
+ @Test
+ void testCreateFlussTableWithExistingEmptyLakeTable() {
+ String database = "empty_lake_db";
+ String tableName = "empty_lake_table";
+ TablePath tablePath = TablePath.of(database, tableName);
+
+ TableDescriptor td = getTableDescriptor(FLUSS_SCHEMA);
+
+ // First creation: prepares the Iceberg table (still empty, no
snapshot)
+ flussIcebergCatalog.createTable(tablePath, td, new
TestingLakeCatalogContext());
+ Table before =
+ flussIcebergCatalog
+ .getIcebergCatalog()
+ .loadTable(TableIdentifier.of(database, tableName));
+ assertThat(before.currentSnapshot()).isNull();
+
+ // A brand-new Fluss table binding to an existing empty compatible
Iceberg table
+ // should be allowed.
+ TestingLakeCatalogContext firstFlussCreationCtx =
+ new TestingLakeCatalogContext() {
+ @Override
+ public boolean isCreatingFlussTable() {
+ return true;
+ }
+ };
+
+ flussIcebergCatalog.createTable(tablePath, td, firstFlussCreationCtx);
+ }
+
+ /**
+ * When creating a fresh Fluss table, if a compatible but non-empty
Iceberg table already
+ * exists, the creation should fail with a clear "not empty" message.
+ */
+ @Test
+ void testCreateFlussTableFailsWhenExistingLakeTableIsNonEmpty() {
+ String database = "non_empty_lake_db";
+ String tableName = "non_empty_lake_table";
+ TablePath tablePath = TablePath.of(database, tableName);
+
+ TableDescriptor td = getTableDescriptor(FLUSS_SCHEMA);
+
+ // First: create the Iceberg table.
+ flussIcebergCatalog.createTable(tablePath, td, new
TestingLakeCatalogContext());
+
+ // Then: append a dummy data file so that the table has a snapshot.
+ Table icebergTable =
+ flussIcebergCatalog
+ .getIcebergCatalog()
+ .loadTable(TableIdentifier.of(database, tableName));
+ appendDummyDataFile(icebergTable);
+ Table afterAppend =
+ flussIcebergCatalog
+ .getIcebergCatalog()
+ .loadTable(TableIdentifier.of(database, tableName));
+ assertThat(afterAppend.currentSnapshot()).isNotNull();
+
+ // Now creating a brand-new Fluss table on top of a non-empty Iceberg
table must fail.
+ TestingLakeCatalogContext firstFlussCreationCtx =
+ new TestingLakeCatalogContext() {
+ @Override
+ public boolean isCreatingFlussTable() {
+ return true;
+ }
+ };
+
+ assertThatThrownBy(
+ () -> flussIcebergCatalog.createTable(tablePath, td,
firstFlussCreationCtx))
+ .isInstanceOf(TableAlreadyExistException.class)
+ .hasMessageContaining("already exists in Iceberg catalog")
+ .hasMessageContaining("not empty");
+ }
+
+ /**
+ * If the existing Iceberg table has an incompatible schema, the creation
should fail regardless
+ * of whether it is a first Fluss creation or a replay.
+ */
+ @Test
+ void testCreateTableFailsWithIncompatibleLakeSchema() {
+ String database = "incompat_lake_db";
+ String tableName = "incompat_lake_table";
+ TablePath tablePath = TablePath.of(database, tableName);
+
+ // First: create the Iceberg table with the base FLUSS_SCHEMA.
+ TableDescriptor tdOrig = getTableDescriptor(FLUSS_SCHEMA);
+ flussIcebergCatalog.createTable(tablePath, tdOrig, new
TestingLakeCatalogContext());
+
+ // Then: attempt to create a Fluss table pointing at the same lake
table, but with a
+ // different schema (an extra column). Should be rejected as
incompatible.
+ Schema mismatchedSchema =
+ Schema.newBuilder()
+ .column("id", DataTypes.BIGINT())
+ .column("name", DataTypes.STRING())
+ .column("amount", DataTypes.INT())
+ .column("address", DataTypes.STRING())
+ .column("extra", DataTypes.STRING())
+ .build();
+ TableDescriptor tdMismatch = getTableDescriptor(mismatchedSchema);
+
+ assertThatThrownBy(
+ () ->
+ flussIcebergCatalog.createTable(
+ tablePath, tdMismatch, new
TestingLakeCatalogContext()))
+ .isInstanceOf(TableAlreadyExistException.class)
+ .hasMessageContaining("already exists in Iceberg catalog")
+ .hasMessageContaining("not compatible");
+ }
+
+ /** Direct unit tests for {@link
IcebergLakeCatalog#isIcebergSchemaCompatibleWithSchema}. */
+ @Test
+ void testIsIcebergSchemaCompatibleWithSchema() {
+ // baseline
+ org.apache.iceberg.Schema base =
+ new org.apache.iceberg.Schema(
+ Types.NestedField.required(1, "id",
Types.LongType.get()),
+ Types.NestedField.optional(2, "name",
Types.StringType.get()));
+
+ // identical schema -> compatible
+ org.apache.iceberg.Schema sameStruct =
+ new org.apache.iceberg.Schema(
+ Types.NestedField.required(11, "id",
Types.LongType.get()),
+ Types.NestedField.optional(22, "name",
Types.StringType.get()));
+
assertThat(flussIcebergCatalog.isIcebergSchemaCompatibleWithSchema(base,
sameStruct))
+ .isTrue();
+
+ // extra column -> incompatible
+ org.apache.iceberg.Schema extra =
+ new org.apache.iceberg.Schema(
+ Types.NestedField.required(1, "id",
Types.LongType.get()),
+ Types.NestedField.optional(2, "name",
Types.StringType.get()),
+ Types.NestedField.optional(3, "extra",
Types.StringType.get()));
+
assertThat(flussIcebergCatalog.isIcebergSchemaCompatibleWithSchema(base,
extra)).isFalse();
+
+ // different name -> incompatible
+ org.apache.iceberg.Schema renamed =
+ new org.apache.iceberg.Schema(
+ Types.NestedField.required(1, "id",
Types.LongType.get()),
+ Types.NestedField.optional(2, "full_name",
Types.StringType.get()));
+
assertThat(flussIcebergCatalog.isIcebergSchemaCompatibleWithSchema(base,
renamed))
+ .isFalse();
+
+ // different type -> incompatible
+ org.apache.iceberg.Schema differentType =
+ new org.apache.iceberg.Schema(
+ Types.NestedField.required(1, "id",
Types.LongType.get()),
+ Types.NestedField.optional(2, "name",
Types.IntegerType.get()));
+
assertThat(flussIcebergCatalog.isIcebergSchemaCompatibleWithSchema(base,
differentType))
+ .isFalse();
+
+ // different nullability -> incompatible
+ org.apache.iceberg.Schema differentNullability =
+ new org.apache.iceberg.Schema(
+ Types.NestedField.required(1, "id",
Types.LongType.get()),
+ Types.NestedField.required(2, "name",
Types.StringType.get()));
+ assertThat(
+
flussIcebergCatalog.isIcebergSchemaCompatibleWithSchema(
+ base, differentNullability))
+ .isFalse();
+ }
+
+ /**
+ * Simulates the real-world case where a user pre-creates the Iceberg
table with a different
+ * bucket count (partition transform {@code bucket[N]}) than what Fluss
expects, then attempts
+ * to bind a Fluss table on top. The mismatch on the partition transform
must be detected and
+ * rejected.
+ */
+ @Test
+ void testCreateTableFailsWithIncompatiblePartitionBucketCount() {
+ String database = "spec_bucket_db";
+ String tableName = "spec_bucket_table";
+ TablePath tablePath = TablePath.of(database, tableName);
+
+ // Pre-create the Iceberg table directly with bucket count = 2.
+ Schema flussSchema =
+ Schema.newBuilder()
+ .column("id", DataTypes.INT())
+ .column("name", DataTypes.STRING())
+ .primaryKey("id")
+ .build();
+ TableDescriptor td2Buckets =
+ TableDescriptor.builder().schema(flussSchema).distributedBy(2,
"id").build();
+ flussIcebergCatalog.createTable(tablePath, td2Buckets, new
TestingLakeCatalogContext());
+
+ // Now Fluss tries to (re-)create the same table with bucket count = 4.
+ TableDescriptor td4Buckets =
+ TableDescriptor.builder().schema(flussSchema).distributedBy(4,
"id").build();
+
+ assertThatThrownBy(
+ () ->
+ flussIcebergCatalog.createTable(
+ tablePath, td4Buckets, new
TestingLakeCatalogContext()))
+ .isInstanceOf(TableAlreadyExistException.class)
+ .hasMessageContaining("already exists in Iceberg catalog")
+ .hasMessageContaining("partition spec is not compatible");
+ }
+
+ /**
+ * When a log table has partition keys, mismatched partition columns
should be rejected as
+ * partition spec incompatible.
+ */
+ @Test
+ void testCreateTableFailsWithIncompatiblePartitionKeys() {
+ String database = "spec_partkeys_db";
+ String tableName = "spec_partkeys_table";
+ TablePath tablePath = TablePath.of(database, tableName);
+
+ Schema flussSchema =
+ Schema.newBuilder()
+ .column("id", DataTypes.BIGINT())
+ .column("name", DataTypes.STRING())
+ .column("region", DataTypes.STRING())
+ .column("dt", DataTypes.STRING())
+ .build();
+
+ // First: partitioned by "region"
+ TableDescriptor tdRegion =
+ TableDescriptor.builder()
+ .schema(flussSchema)
+ .distributedBy(3)
+ .partitionedBy("region")
+ .build();
+ flussIcebergCatalog.createTable(tablePath, tdRegion, new
TestingLakeCatalogContext());
+
+ // Second: partitioned by "dt" -> mismatched
+ TableDescriptor tdDt =
+ TableDescriptor.builder()
+ .schema(flussSchema)
+ .distributedBy(3)
+ .partitionedBy("dt")
+ .build();
+ assertThatThrownBy(
+ () ->
+ flussIcebergCatalog.createTable(
+ tablePath, tdDt, new
TestingLakeCatalogContext()))
+ .isInstanceOf(TableAlreadyExistException.class)
+ .hasMessageContaining("partition spec is not compatible");
+ }
+
+ /**
+ * If someone alters the sort order of the existing Iceberg table (e.g.
removes the fluss sort
+ * order), Fluss should detect the mismatch when re-binding.
+ */
+ @Test
+ void testCreateTableFailsWithIncompatibleSortOrder() {
+ String database = "sortorder_db";
+ String tableName = "sortorder_table";
+ TablePath tablePath = TablePath.of(database, tableName);
+
+ TableDescriptor td = getTableDescriptor(FLUSS_SCHEMA);
+ flussIcebergCatalog.createTable(tablePath, td, new
TestingLakeCatalogContext());
+
+ // Manually replace the sort order of the existing Iceberg table.
+ Table iceberg =
+ flussIcebergCatalog
+ .getIcebergCatalog()
+ .loadTable(TableIdentifier.of(database, tableName));
+ iceberg.replaceSortOrder().asc("id").commit();
+
+ assertThatThrownBy(
+ () ->
+ flussIcebergCatalog.createTable(
+ tablePath, td, new
TestingLakeCatalogContext()))
+ .isInstanceOf(TableAlreadyExistException.class)
+ .hasMessageContaining("sort order is not compatible")
+ .hasMessageContaining("ASC(__offset)")
+ .hasMessageContaining("pre-create the Iceberg table without a
custom sort order");
+ }
+
+ /**
+ * If someone changes a fluss-managed property on the existing Iceberg
table (e.g. write mode),
+ * the mismatch should be surfaced.
+ */
+ @Test
+ void testCreateTableFailsWithIncompatibleProperties() {
+ String database = "props_db";
+ String tableName = "props_table";
+ TablePath tablePath = TablePath.of(database, tableName);
+
+ // pk table -> we set MERGE_MODE / UPDATE_MODE / DELETE_MODE to
merge-on-read.
+ Schema pkSchema =
+ Schema.newBuilder()
+ .column("id", DataTypes.INT())
+ .column("name", DataTypes.STRING())
+ .primaryKey("id")
+ .build();
+ TableDescriptor td =
+ TableDescriptor.builder().schema(pkSchema).distributedBy(4,
"id").build();
+ flussIcebergCatalog.createTable(tablePath, td, new
TestingLakeCatalogContext());
+
+ // Manually change a required property to an incompatible value.
+ Table iceberg =
+ flussIcebergCatalog
+ .getIcebergCatalog()
+ .loadTable(TableIdentifier.of(database, tableName));
+ iceberg.updateProperties()
+ .set(TableProperties.MERGE_MODE,
RowLevelOperationMode.COPY_ON_WRITE.modeName())
+ .commit();
+
+ assertThatThrownBy(
+ () ->
+ flussIcebergCatalog.createTable(
+ tablePath, td, new
TestingLakeCatalogContext()))
+ .isInstanceOf(TableAlreadyExistException.class)
+ .hasMessageContaining("table properties are not compatible");
+ }
+
+ /** White-box tests for {@link
IcebergLakeCatalog#isIcebergPartitionSpecCompatible}. */
+ @Test
+ void testIsIcebergPartitionSpecCompatible() {
+ // baseline schema with an "id" column (required).
+ org.apache.iceberg.Schema schemaA =
+ new org.apache.iceberg.Schema(
+ Types.NestedField.required(1, "id",
Types.LongType.get()),
+ Types.NestedField.optional(2, "name",
Types.StringType.get()));
+ org.apache.iceberg.Schema schemaB =
+ new org.apache.iceberg.Schema(
+ Types.NestedField.required(11, "id",
Types.LongType.get()),
+ Types.NestedField.optional(22, "name",
Types.StringType.get()));
+
+ PartitionSpec bucket4A =
PartitionSpec.builderFor(schemaA).bucket("id", 4).build();
+ PartitionSpec bucket4B =
PartitionSpec.builderFor(schemaB).bucket("id", 4).build();
+ PartitionSpec bucket2A =
PartitionSpec.builderFor(schemaA).bucket("id", 2).build();
+ PartitionSpec identityNameA =
PartitionSpec.builderFor(schemaA).identity("name").build();
+
+ // Same transform + same source column, different field ids ->
compatible.
+ assertThat(
+ flussIcebergCatalog.isIcebergPartitionSpecCompatible(
+ bucket4A, schemaA, bucket4B, schemaB))
+ .isTrue();
+
+ // Different bucket count on same column -> incompatible.
+ assertThat(
+ flussIcebergCatalog.isIcebergPartitionSpecCompatible(
+ bucket4A, schemaA, bucket2A, schemaA))
+ .isFalse();
+
+ // Different source column -> incompatible.
+ assertThat(
+ flussIcebergCatalog.isIcebergPartitionSpecCompatible(
+ bucket4A, schemaA, identityNameA, schemaA))
+ .isFalse();
+
+ // Different partition field count -> incompatible.
+ PartitionSpec twoFields =
+
PartitionSpec.builderFor(schemaA).identity("name").bucket("id", 4).build();
+ assertThat(
+ flussIcebergCatalog.isIcebergPartitionSpecCompatible(
+ bucket4A, schemaA, twoFields, schemaA))
+ .isFalse();
+ }
+
+ /** White-box tests for {@link
IcebergLakeCatalog#isIcebergSortOrderCompatible}. */
+ @Test
+ void testIsIcebergSortOrderCompatible() {
+ org.apache.iceberg.Schema schemaA =
+ new org.apache.iceberg.Schema(
+ Types.NestedField.required(1, "id",
Types.LongType.get()),
+ Types.NestedField.optional(2, "name",
Types.StringType.get()));
+ org.apache.iceberg.Schema schemaB =
+ new org.apache.iceberg.Schema(
+ Types.NestedField.required(11, "id",
Types.LongType.get()),
+ Types.NestedField.optional(22, "name",
Types.StringType.get()));
+
+ SortOrder ascIdA = SortOrder.builderFor(schemaA).asc("id").build();
+ SortOrder ascIdB = SortOrder.builderFor(schemaB).asc("id").build();
+ SortOrder descIdA = SortOrder.builderFor(schemaA).desc("id").build();
+ SortOrder ascNameA = SortOrder.builderFor(schemaA).asc("name").build();
+
+ // Same source + direction, different ids -> compatible.
+ assertThat(
+ flussIcebergCatalog.isIcebergSortOrderCompatible(
+ ascIdA, schemaA, ascIdB, schemaB))
+ .isTrue();
+
+ // Different direction -> incompatible.
+ assertThat(
+ flussIcebergCatalog.isIcebergSortOrderCompatible(
+ ascIdA, schemaA, descIdA, schemaA))
+ .isFalse();
+
+ // Different source column -> incompatible.
+ assertThat(
+ flussIcebergCatalog.isIcebergSortOrderCompatible(
+ ascIdA, schemaA, ascNameA, schemaA))
+ .isFalse();
+
+ // Different sort field count -> incompatible.
+ SortOrder twoFields =
SortOrder.builderFor(schemaA).asc("id").asc("name").build();
+ assertThat(
+ flussIcebergCatalog.isIcebergSortOrderCompatible(
+ ascIdA, schemaA, twoFields, schemaA))
+ .isFalse();
+ }
+
+ /** White-box tests for {@link
IcebergLakeCatalog#isIcebergPropertiesCompatible}. */
+ @Test
+ void testIsIcebergPropertiesCompatible() {
+ Map<String, String> expected = new java.util.HashMap<>();
+ expected.put("k1", "v1");
+ expected.put("k2", "v2");
+
+ // Existing has all expected keys with same values -> compatible.
+ Map<String, String> existing = new java.util.HashMap<>(expected);
+ assertThat(flussIcebergCatalog.isIcebergPropertiesCompatible(existing,
expected)).isTrue();
+
+ // Existing also has extra keys -> still compatible (extras are
ignored).
+ existing = new java.util.HashMap<>(expected);
+ existing.put("engine.internal", "whatever");
+ assertThat(flussIcebergCatalog.isIcebergPropertiesCompatible(existing,
expected)).isTrue();
+
+ // Existing missing an expected key -> incompatible.
+ existing = new java.util.HashMap<>();
+ existing.put("k1", "v1");
+ assertThat(flussIcebergCatalog.isIcebergPropertiesCompatible(existing,
expected)).isFalse();
+
+ // Existing has key but different value -> incompatible.
+ existing = new java.util.HashMap<>(expected);
+ existing.put("k2", "different");
+ assertThat(flussIcebergCatalog.isIcebergPropertiesCompatible(existing,
expected)).isFalse();
+ }
+
+ private void appendDummyDataFile(Table icebergTable) {
+ DataFile dataFile =
+ DataFiles.builder(PartitionSpec.unpartitioned())
+ .withPath("/tmp/fluss-fake-file.parquet")
+ .withFileSizeInBytes(1024)
+ .withRecordCount(1L)
+ .withFormat(FileFormat.PARQUET)
+ .build();
+ icebergTable.newAppend().appendFile(dataFile).commit();
+ }
Review Comment:
`appendDummyDataFile` builds the `DataFile` using
`PartitionSpec.unpartitioned()`, but Fluss-created Iceberg tables always have a
non-empty partition spec (at least `identity(__bucket)`). Appending a file with
an unpartitioned spec is incompatible and will likely fail/leave the table
without a snapshot, making
`testCreateFlussTableFailsWhenExistingLakeTableIsNonEmpty` unreliable.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]