diff --git a/.github/workflows/pr_build_linux.yml b/.github/workflows/pr_build_linux.yml index 6b7130b401..af0aa75934 100644 --- a/.github/workflows/pr_build_linux.yml +++ b/.github/workflows/pr_build_linux.yml @@ -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 diff --git a/.github/workflows/pr_build_macos.yml b/.github/workflows/pr_build_macos.yml index 65d87ec24f..531cd210b0 100644 --- a/.github/workflows/pr_build_macos.yml +++ b/.github/workflows/pr_build_macos.yml @@ -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 diff --git a/spark/src/main/scala/org/apache/comet/serde/operator/CometIcebergNativeScan.scala b/spark/src/main/scala/org/apache/comet/serde/operator/CometIcebergNativeScan.scala index 0875ea04a8..cb8cc5752b 100644 --- a/spark/src/main/scala/org/apache/comet/serde/operator/CometIcebergNativeScan.scala +++ b/spark/src/main/scala/org/apache/comet/serde/operator/CometIcebergNativeScan.scala @@ -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 = @@ -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. * diff --git a/spark/src/test/scala/org/apache/comet/serde/operator/CometIcebergDeleteFileSerdeSuite.scala b/spark/src/test/scala/org/apache/comet/serde/operator/CometIcebergDeleteFileSerdeSuite.scala new file mode 100644 index 0000000000..5b3d6b5558 --- /dev/null +++ b/spark/src/test/scala/org/apache/comet/serde/operator/CometIcebergDeleteFileSerdeSuite.scala @@ -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 + } +}