Skip to content
Open
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
1 change: 1 addition & 0 deletions core/src/main/java/org/apache/iceberg/avro/Avro.java
Original file line number Diff line number Diff line change
Expand Up @@ -506,6 +506,7 @@ public <T> EqualityDeleteWriter<T> buildEqualityWriter() throws IOException {
Preconditions.checkArgument(
spec.isUnpartitioned() || partition != null,
"Partition must not be null for partitioned writes");
DeleteSchemaUtil.validateEqualityFieldIds(equalityFieldIds, rowSchema);

meta("delete-type", "equality");
meta(
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -36,6 +36,7 @@
import org.apache.iceberg.encryption.EncryptedOutputFile;
import org.apache.iceberg.encryption.EncryptionKeyMetadata;
import org.apache.iceberg.io.DataWriter;
import org.apache.iceberg.io.DeleteSchemaUtil;
import org.apache.iceberg.io.FileWriter;
import org.apache.iceberg.relocated.com.google.common.base.Preconditions;

Expand Down Expand Up @@ -213,6 +214,9 @@ protected void validate() {
spec.isUnpartitioned() || partition != null,
"Invalid partition, does not match spec: %s",
spec);
if (content == FileContent.EQUALITY_DELETES) {
DeleteSchemaUtil.validateEqualityFieldIds(equalityFieldIds, schema);
}
}

/** Builder for creating {@link DataWriter} instances for writing data files. */
Expand Down
20 changes: 20 additions & 0 deletions core/src/main/java/org/apache/iceberg/io/DeleteSchemaUtil.java
Original file line number Diff line number Diff line change
Expand Up @@ -18,8 +18,11 @@
*/
package org.apache.iceberg.io;

import java.util.Set;
import org.apache.iceberg.MetadataColumns;
import org.apache.iceberg.Schema;
import org.apache.iceberg.relocated.com.google.common.base.Preconditions;
import org.apache.iceberg.relocated.com.google.common.collect.Sets;
import org.apache.iceberg.types.Types;

public class DeleteSchemaUtil {
Expand Down Expand Up @@ -54,4 +57,21 @@ public static Schema posDeleteReadSchema(Schema rowSchema) {
rowSchema.asStruct(),
MetadataColumns.DELETE_FILE_ROW_DOC));
}

public static void validateEqualityFieldIds(int[] equalityFieldIds, Schema equalityDeleteSchema) {
Preconditions.checkArgument(
equalityFieldIds != null && equalityFieldIds.length > 0,
"Equality delete field IDs must not be null or empty");
Preconditions.checkArgument(equalityDeleteSchema != null, "Schema must not be null");

Set<Integer> seenFieldIds = Sets.newHashSetWithExpectedSize(equalityFieldIds.length);
for (int fieldId : equalityFieldIds) {
Preconditions.checkArgument(
seenFieldIds.add(fieldId), "Duplicate equality delete field ID: %s", fieldId);
Preconditions.checkArgument(
equalityDeleteSchema.findField(fieldId) != null,
"Invalid equality delete field ID: %s",
fieldId);
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@
package org.apache.iceberg.avro;

import static org.assertj.core.api.Assertions.assertThat;
import static org.assertj.core.api.Assertions.assertThatThrownBy;

import java.io.File;
import java.io.IOException;
Expand Down Expand Up @@ -112,6 +113,57 @@ public void testEqualityDeleteWriter() throws IOException {
assertThat(deletedRecords).as("Deleted records should match expected").isEqualTo(records);
}

@Test
public void equalityDeleteWriterRejectsEmptyEqualityFieldIds() {
OutputFile out = new InMemoryOutputFile();

assertThatThrownBy(
() ->
Avro.writeDeletes(out)
.createWriterFunc(DataWriter::create)
.overwrite()
.rowSchema(SCHEMA)
.withSpec(PartitionSpec.unpartitioned())
.equalityFieldIds(new int[0])
.buildEqualityWriter())
.isInstanceOf(IllegalArgumentException.class)
.hasMessageContaining("Equality delete field IDs must not be null or empty");
}

@Test
public void equalityDeleteWriterRejectsDuplicateEqualityFieldIds() {
OutputFile out = new InMemoryOutputFile();

assertThatThrownBy(
() ->
Avro.writeDeletes(out)
.createWriterFunc(DataWriter::create)
.overwrite()
.rowSchema(SCHEMA)
.withSpec(PartitionSpec.unpartitioned())
.equalityFieldIds(1, 1)
.buildEqualityWriter())
.isInstanceOf(IllegalArgumentException.class)
.hasMessageContaining("Duplicate equality delete field ID: 1");
}

@Test
public void equalityDeleteWriterRejectsMissingEqualityFieldId() {
OutputFile out = new InMemoryOutputFile();

assertThatThrownBy(
() ->
Avro.writeDeletes(out)
.createWriterFunc(DataWriter::create)
.overwrite()
.rowSchema(SCHEMA)
.withSpec(PartitionSpec.unpartitioned())
.equalityFieldIds(99)
.buildEqualityWriter())
.isInstanceOf(IllegalArgumentException.class)
.hasMessageContaining("Invalid equality delete field ID: 99");
}

@Test
public void testPositionDeleteWriter() throws IOException {
Schema deleteSchema =
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -23,9 +23,20 @@
import static org.assertj.core.api.AssertionsForInterfaceTypes.assertThat;

import java.lang.reflect.Method;
import java.nio.ByteBuffer;
import org.apache.iceberg.FileContent;
import org.apache.iceberg.FileFormat;
import org.apache.iceberg.Metrics;
import org.apache.iceberg.MetricsConfig;
import org.apache.iceberg.PartitionSpec;
import org.apache.iceberg.Schema;
import org.apache.iceberg.encryption.EncryptedFiles;
import org.apache.iceberg.encryption.EncryptedOutputFile;
import org.apache.iceberg.encryption.EncryptionKeyMetadata;
import org.apache.iceberg.inmemory.InMemoryOutputFile;
import org.apache.iceberg.io.FileAppender;
import org.apache.iceberg.io.InputFile;
import org.apache.iceberg.types.Types;
import org.apache.iceberg.util.Pair;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
Expand Down Expand Up @@ -124,6 +135,25 @@ public static void register() {
}
}

@Test
void equalityDeleteWriterRejectsMissingEqualityFieldId() {
FormatModelRegistry.register(new DummyParquetFormatModel(Object.class, Object.class));
EncryptedOutputFile outputFile =
EncryptedFiles.encryptedOutput(new InMemoryOutputFile(), EncryptionKeyMetadata.EMPTY);
Schema schema = new Schema(Types.NestedField.required(1, "id", Types.LongType.get()));

assertThatThrownBy(
() ->
FormatModelRegistry.equalityDeleteWriteBuilder(
FileFormat.PARQUET, Object.class, outputFile)
.schema(schema)
.spec(PartitionSpec.unpartitioned())
.equalityFieldIds(99)
.build())
.isInstanceOf(IllegalArgumentException.class)
.hasMessageContaining("Invalid equality delete field ID: 99");
}

private static class DummyParquetFormatModel implements FormatModel<Object, Object> {
private final Class<?> type;
private final Class<?> schemaType;
Expand Down Expand Up @@ -152,12 +182,82 @@ public Class<Object> schemaType() {

@Override
public ModelWriteBuilder<Object, Object> writeBuilder(EncryptedOutputFile outputFile) {
return null;
return new DummyModelWriteBuilder();
}

@Override
public ReadBuilder<Object, Object> readBuilder(InputFile inputFile) {
return null;
}
}

private static class DummyModelWriteBuilder implements ModelWriteBuilder<Object, Object> {
@Override
public ModelWriteBuilder<Object, Object> schema(Schema schema) {
return this;
}

@Override
public ModelWriteBuilder<Object, Object> engineSchema(Object schema) {
return this;
}

@Override
public ModelWriteBuilder<Object, Object> set(String property, String value) {
return this;
}

@Override
public ModelWriteBuilder<Object, Object> meta(String property, String value) {
return this;
}

@Override
public ModelWriteBuilder<Object, Object> content(FileContent content) {
return this;
}

@Override
public ModelWriteBuilder<Object, Object> metricsConfig(MetricsConfig metricsConfig) {
return this;
}

@Override
public ModelWriteBuilder<Object, Object> overwrite() {
return this;
}

@Override
public ModelWriteBuilder<Object, Object> withFileEncryptionKey(ByteBuffer encryptionKey) {
return this;
}

@Override
public ModelWriteBuilder<Object, Object> withAADPrefix(ByteBuffer aadPrefix) {
return this;
}

@Override
public FileAppender<Object> build() {
return new NoOpFileAppender();
}
}

private static class NoOpFileAppender implements FileAppender<Object> {
@Override
public void add(Object datum) {}

@Override
public Metrics metrics() {
return null;
}

@Override
public long length() {
return 0;
}

@Override
public void close() {}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -35,6 +35,7 @@
import org.apache.iceberg.encryption.EncryptedOutputFile;
import org.apache.iceberg.encryption.EncryptionKeyMetadata;
import org.apache.iceberg.io.DataWriter;
import org.apache.iceberg.io.DeleteSchemaUtil;
import org.apache.iceberg.io.FileWriterFactory;
import org.apache.iceberg.orc.ORC;
import org.apache.iceberg.parquet.Parquet;
Expand Down Expand Up @@ -69,6 +70,10 @@ protected BaseFileWriterFactory(
Schema equalityDeleteRowSchema,
SortOrder equalityDeleteSortOrder,
Map<String, String> writerProperties) {
if (equalityDeleteRowSchema != null && equalityFieldIds != null) {
DeleteSchemaUtil.validateEqualityFieldIds(equalityFieldIds, equalityDeleteRowSchema);
}

this.table = table;
this.dataFileFormat = dataFileFormat;
this.dataSchema = dataSchema;
Expand All @@ -92,6 +97,10 @@ protected BaseFileWriterFactory(
SortOrder equalityDeleteSortOrder,
Schema positionDeleteRowSchema,
Map<String, String> writerProperties) {
if (equalityDeleteRowSchema != null && equalityFieldIds != null) {
DeleteSchemaUtil.validateEqualityFieldIds(equalityFieldIds, equalityDeleteRowSchema);
}

this.table = table;
this.dataFileFormat = dataFileFormat;
this.dataSchema = dataSchema;
Expand All @@ -115,6 +124,10 @@ protected BaseFileWriterFactory(
Schema equalityDeleteRowSchema,
SortOrder equalityDeleteSortOrder,
Schema positionDeleteRowSchema) {
if (equalityDeleteRowSchema != null && equalityFieldIds != null) {
DeleteSchemaUtil.validateEqualityFieldIds(equalityFieldIds, equalityDeleteRowSchema);
}

this.table = table;
this.dataFileFormat = dataFileFormat;
this.dataSchema = dataSchema;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -35,6 +35,7 @@
import org.apache.iceberg.deletes.PositionDeleteWriter;
import org.apache.iceberg.encryption.EncryptedOutputFile;
import org.apache.iceberg.encryption.EncryptionUtil;
import org.apache.iceberg.io.DeleteSchemaUtil;
import org.apache.iceberg.io.FileAppender;
import org.apache.iceberg.io.FileAppenderFactory;
import org.apache.iceberg.io.OutputFile;
Expand Down Expand Up @@ -132,6 +133,10 @@ public GenericAppenderFactory(
int[] equalityFieldIds,
Schema eqDeleteRowSchema,
Schema posDeleteRowSchema) {
if (eqDeleteRowSchema != null && equalityFieldIds != null) {
DeleteSchemaUtil.validateEqualityFieldIds(equalityFieldIds, eqDeleteRowSchema);
}

this.table = table;
this.config = config == null ? Maps.newHashMap() : config;

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -36,6 +36,7 @@
import org.apache.iceberg.formats.FileWriterBuilder;
import org.apache.iceberg.formats.FormatModelRegistry;
import org.apache.iceberg.io.DataWriter;
import org.apache.iceberg.io.DeleteSchemaUtil;
import org.apache.iceberg.io.FileWriterFactory;
import org.apache.iceberg.relocated.com.google.common.base.Preconditions;
import org.apache.iceberg.relocated.com.google.common.collect.ImmutableMap;
Expand Down Expand Up @@ -73,6 +74,10 @@ protected RegistryBasedFileWriterFactory(
Map<String, String> writerProperties,
S inputSchema,
S equalityDeleteInputSchema) {
if (equalityDeleteRowSchema != null && equalityFieldIds != null) {
DeleteSchemaUtil.validateEqualityFieldIds(equalityFieldIds, equalityDeleteRowSchema);
}

this.table = table;
this.dataFileFormat = dataFileFormat;
this.inputType = inputType;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -136,4 +136,33 @@ void createFactoryWithConflictConfig() {
.hasMessageContaining(
"Cannot set metrics properties when the table is provided, use table properties instead");
}

@TestTemplate
void equalityFieldIdsAreValidatedAgainstEqualityDeleteRowSchema() {
int equalityFieldId = table.schema().findField("id").fieldId();

assertThatNoException()
.isThrownBy(
() ->
new GenericAppenderFactory(
null,
PartitionSpec.unpartitioned(),
new int[] {equalityFieldId},
table.schema().select("id")));
}

@TestTemplate
void equalityFieldIdsAreRejectedWhenMissingFromEqualityDeleteRowSchema() {
int equalityFieldId = table.schema().findField("id").fieldId();

assertThatThrownBy(
() ->
new GenericAppenderFactory(
table.schema().select("id"),
PartitionSpec.unpartitioned(),
new int[] {equalityFieldId},
table.schema().select("data")))
.isInstanceOf(IllegalArgumentException.class)
.hasMessage("Invalid equality delete field ID: %s", equalityFieldId);
}
}
Loading
Loading