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 .github/workflows/pr_build_linux.yml
Original file line number Diff line number Diff line change
Expand Up @@ -317,6 +317,7 @@ jobs:
org.apache.comet.CometIcebergWriteActionSuite
org.apache.comet.CometIcebergWriteDetectionSuite
org.apache.comet.iceberg.IcebergReflectionSuite
org.apache.comet.serde.operator.CometIcebergDeleteFileSerdeSuite
org.apache.comet.csv.CometCsvNativeReadSuite
org.apache.comet.CometFuzzTestSuite
org.apache.comet.CometFuzzIcebergSuite
Expand Down
1 change: 1 addition & 0 deletions .github/workflows/pr_build_macos.yml
Original file line number Diff line number Diff line change
Expand Up @@ -133,6 +133,7 @@ jobs:
org.apache.comet.CometIcebergWriteActionSuite
org.apache.comet.CometIcebergWriteDetectionSuite
org.apache.comet.iceberg.IcebergReflectionSuite
org.apache.comet.serde.operator.CometIcebergDeleteFileSerdeSuite
org.apache.comet.csv.CometCsvNativeReadSuite
org.apache.comet.CometFuzzTestSuite
org.apache.comet.CometFuzzIcebergSuite
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -280,62 +280,9 @@ object CometIcebergNativeScan extends CometOperatorSerde[CometBatchScanExec] wit

val deletes = IcebergReflection.getDeleteFilesFromTask(task, fileScanTaskClass)

deletes.asScala.map { deleteFile =>
// The path is the one essential field. A delete file we cannot locate cannot be applied,
// and silently skipping it would leak deleted rows, so treat a missing path as fatal.
val deletePath = IcebergReflection
.extractFileLocation(contentFileClass, deleteFile)
.getOrElse(
throw new RuntimeException(
"Neither location() nor path() is declared on this Iceberg version's " +
"ContentFile -- cannot extract delete file path from FileScanTask"))

val deleteBuilder = OperatorOuterClass.IcebergDeleteFile.newBuilder()
deleteBuilder.setFilePath(deletePath)

val contentType =
try {
val contentMethod = IcebergReflection.getMethod(deleteFileClass, "content")
val content = contentMethod.invoke(deleteFile)
content.toString match {
case IcebergReflection.ContentTypes.POSITION_DELETES =>
IcebergReflection.ContentTypes.POSITION_DELETES
case IcebergReflection.ContentTypes.EQUALITY_DELETES =>
IcebergReflection.ContentTypes.EQUALITY_DELETES
case other => other
}
} catch {
case _: Exception =>
IcebergReflection.ContentTypes.POSITION_DELETES
}
deleteBuilder.setContentType(contentType)

val specId =
try {
val specIdMethod = IcebergReflection.getMethod(deleteFileClass, "specId")
specIdMethod.invoke(deleteFile).asInstanceOf[Int]
} catch {
case _: Exception => 0
}
deleteBuilder.setPartitionSpecId(specId)

try {
val equalityIdsMethod =
IcebergReflection.getMethod(deleteFileClass, "equalityFieldIds")
val equalityIds = equalityIdsMethod
.invoke(deleteFile)
.asInstanceOf[java.util.List[Integer]]
equalityIds.forEach(id => deleteBuilder.addEqualityIds(id))
} catch {
case _: Exception =>
}

// Encrypted delete files carry a plaintext StandardKeyMetadata blob; forward it verbatim.
// Unencrypted delete files leave the field unset.
keyMetadataBytes(keyMetadataMethod, deleteFile).foreach(deleteBuilder.setKeyMetadata)

deleteBuilder.build()
}.toSeq
deletes.asScala
.map(serializeDeleteFile(_, contentFileClass, deleteFileClass, keyMetadataMethod))
.toSeq
} catch {
case e: Exception =>
val msg =
Expand All @@ -346,6 +293,63 @@ object CometIcebergNativeScan extends CometOperatorSerde[CometBatchScanExec] wit
}
}

/**
* Serializes a single Iceberg DeleteFile to protobuf.
*
* `content()`, `specId()`, and `equalityFieldIds()` are declared on the public `ContentFile` /
* `DeleteFile` interfaces across all supported Iceberg versions, so a `getMethod` miss or an
* `invoke` failure on any of them means something is genuinely wrong. None of the three may
* fall back to a default: the scan is already committed to native execution, and a guessed
* content type, partition spec, or dropped equality keys all silently return wrong rows.
* Failures propagate to `extractDeleteFilesList`'s outer catch.
*/
private[operator] def serializeDeleteFile(
deleteFile: Any,
contentFileClass: Class[_],
deleteFileClass: Class[_],
keyMetadataMethod: Method): OperatorOuterClass.IcebergDeleteFile = {
// The path is the one essential field. A delete file we cannot locate cannot be applied,
// and silently skipping it would leak deleted rows, so treat a missing path as fatal.
val deletePath = IcebergReflection
.extractFileLocation(contentFileClass, deleteFile)
.getOrElse(
throw new RuntimeException(
"Neither location() nor path() is declared on this Iceberg version's " +
"ContentFile -- cannot extract delete file path from FileScanTask"))

val deleteBuilder = OperatorOuterClass.IcebergDeleteFile.newBuilder()
deleteBuilder.setFilePath(deletePath)

val contentMethod = IcebergReflection.getMethod(deleteFileClass, "content")
val contentType = contentMethod.invoke(deleteFile).toString match {
case IcebergReflection.ContentTypes.POSITION_DELETES =>
IcebergReflection.ContentTypes.POSITION_DELETES
case IcebergReflection.ContentTypes.EQUALITY_DELETES =>
IcebergReflection.ContentTypes.EQUALITY_DELETES
case other => other
}
deleteBuilder.setContentType(contentType)

val specIdMethod = IcebergReflection.getMethod(deleteFileClass, "specId")
deleteBuilder.setPartitionSpecId(specIdMethod.invoke(deleteFile).asInstanceOf[Int])

val equalityIdsMethod = IcebergReflection.getMethod(deleteFileClass, "equalityFieldIds")
val equalityIds = equalityIdsMethod.invoke(deleteFile).asInstanceOf[java.util.List[Integer]]
// Iceberg's BaseFile stores equality field IDs in a nullable backing array, so
// equalityFieldIds() returns null for files without equality keys. A null return is a
// normal accessor result, unlike a reflective lookup or invocation failure, and does not
// make serialization fail.
if (equalityIds != null) {
equalityIds.forEach(id => deleteBuilder.addEqualityIds(id))
}

// Encrypted delete files carry a plaintext StandardKeyMetadata blob; forward it verbatim.
// Unencrypted delete files leave the field unset.
keyMetadataBytes(keyMetadataMethod, deleteFile).foreach(deleteBuilder.setKeyMetadata)

deleteBuilder.build()
}

/**
* Serializes partition spec and data from an Iceberg FileScanTask.
*
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,130 @@
/*
* 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.comet.serde.operator

import java.lang.reflect.InvocationTargetException

import org.scalatest.funsuite.AnyFunSuite

/**
* Locks in the fail-loud behavior of [[CometIcebergNativeScan.serializeDeleteFile]] required by
* apache/datafusion-comet#5256: on supported Iceberg versions `content()`, `specId()`, and
* `equalityFieldIds()` are always declared, so a reflective invocation failure must propagate
* rather than fall back to a guessed value. A null `equalityFieldIds()` (a position-delete file)
* stays a legitimate "no equality keys" result.
*/
class CometIcebergDeleteFileSerdeSuite extends AnyFunSuite {

private def keyMetadataMethod(clazz: Class[_]) = clazz.getMethod("keyMetadata")

private def serialize(file: AnyRef) =
CometIcebergNativeScan.serializeDeleteFile(
file,
file.getClass,
file.getClass,
keyMetadataMethod(file.getClass))

test("position-delete file: null equalityFieldIds() serializes with no equality ids") {
val proto = serialize(new PositionDeleteFile)
assert(proto.getContentType == "POSITION_DELETES")
assert(proto.getPartitionSpecId == 7)
assert(proto.getEqualityIdsCount == 0)
assert(proto.getFilePath == "s3://bucket/pos-delete.parquet")
}

test("equality-delete file: declared equalityFieldIds() are serialized") {
val proto = serialize(new EqualityDeleteFile)
assert(proto.getContentType == "EQUALITY_DELETES")
assert(proto.getEqualityIdsCount == 2)
assert(proto.getEqualityIds(0) == 3)
assert(proto.getEqualityIds(1) == 5)
}

test("content() invocation failure propagates instead of defaulting to POSITION_DELETES") {
val ex = intercept[InvocationTargetException](serialize(new ThrowingContentDeleteFile))
assert(ex.getCause.getMessage == "content boom")
}

test("specId() invocation failure propagates instead of defaulting to 0") {
val ex = intercept[InvocationTargetException](serialize(new ThrowingSpecIdDeleteFile))
assert(ex.getCause.getMessage == "spec boom")
}

test("equalityFieldIds() invocation failure propagates instead of dropping equality ids") {
val ex = intercept[InvocationTargetException](serialize(new ThrowingEqualityIdsDeleteFile))
assert(ex.getCause.getMessage == "ids boom")
}

test("missing content() accessor is fatal, not a default") {
// getMethod throws NoSuchMethodException directly (not wrapped) when the accessor is absent.
assertThrows[NoSuchMethodException](serialize(new NoContentAccessorDeleteFile))
}

// -- Synthetic DeleteFile stubs. Each declares the full accessor set serializeDeleteFile
// resolves; only the field under test misbehaves. --

class PositionDeleteFile {
def location(): String = "s3://bucket/pos-delete.parquet"
def content(): String = "POSITION_DELETES"
def specId(): Int = 7
def equalityFieldIds(): java.util.List[Integer] = null
def keyMetadata(): java.nio.ByteBuffer = null
}

class EqualityDeleteFile {
def location(): String = "s3://bucket/eq-delete.parquet"
def content(): String = "EQUALITY_DELETES"
def specId(): Int = 0
def equalityFieldIds(): java.util.List[Integer] =
java.util.List.of(Integer.valueOf(3), Integer.valueOf(5))
def keyMetadata(): java.nio.ByteBuffer = null
}

class ThrowingContentDeleteFile {
def location(): String = "s3://bucket/d.parquet"
def content(): String = throw new RuntimeException("content boom")
def specId(): Int = 0
def equalityFieldIds(): java.util.List[Integer] = null
def keyMetadata(): java.nio.ByteBuffer = null
}

class ThrowingSpecIdDeleteFile {
def location(): String = "s3://bucket/d.parquet"
def content(): String = "POSITION_DELETES"
def specId(): Int = throw new RuntimeException("spec boom")
def equalityFieldIds(): java.util.List[Integer] = null
def keyMetadata(): java.nio.ByteBuffer = null
}

class ThrowingEqualityIdsDeleteFile {
def location(): String = "s3://bucket/d.parquet"
def content(): String = "EQUALITY_DELETES"
def specId(): Int = 0
def equalityFieldIds(): java.util.List[Integer] = throw new RuntimeException("ids boom")
def keyMetadata(): java.nio.ByteBuffer = null
}

class NoContentAccessorDeleteFile {
def location(): String = "s3://bucket/d.parquet"
def specId(): Int = 0
def equalityFieldIds(): java.util.List[Integer] = null
def keyMetadata(): java.nio.ByteBuffer = null
}
}
Loading