diff --git a/.github/trigger_files/IO_Iceberg_Integration_Tests.json b/.github/trigger_files/IO_Iceberg_Integration_Tests.json index 34a6e02150e7..b73af5e61a43 100644 --- a/.github/trigger_files/IO_Iceberg_Integration_Tests.json +++ b/.github/trigger_files/IO_Iceberg_Integration_Tests.json @@ -1,4 +1,4 @@ { "comment": "Modify this file in a trivial way to cause this test suite to run.", - "modification": 4 + "modification": 1 } diff --git a/.github/trigger_files/IO_Iceberg_Integration_Tests_Dataflow.json b/.github/trigger_files/IO_Iceberg_Integration_Tests_Dataflow.json index 5abe02fc09c7..3a009261f4f9 100644 --- a/.github/trigger_files/IO_Iceberg_Integration_Tests_Dataflow.json +++ b/.github/trigger_files/IO_Iceberg_Integration_Tests_Dataflow.json @@ -1,4 +1,4 @@ { "comment": "Modify this file in a trivial way to cause this test suite to run.", - "modification": 1 + "modification": 2 } diff --git a/.github/trigger_files/IO_Iceberg_Managed_Integration_Tests_Dataflow.json b/.github/trigger_files/IO_Iceberg_Managed_Integration_Tests_Dataflow.json index 5abe02fc09c7..3a009261f4f9 100644 --- a/.github/trigger_files/IO_Iceberg_Managed_Integration_Tests_Dataflow.json +++ b/.github/trigger_files/IO_Iceberg_Managed_Integration_Tests_Dataflow.json @@ -1,4 +1,4 @@ { "comment": "Modify this file in a trivial way to cause this test suite to run.", - "modification": 1 + "modification": 2 } diff --git a/.github/trigger_files/beam_PostCommit_Python_Xlang_IO_Direct.json b/.github/trigger_files/beam_PostCommit_Python_Xlang_IO_Direct.json index e3d6056a5de9..b26833333238 100644 --- a/.github/trigger_files/beam_PostCommit_Python_Xlang_IO_Direct.json +++ b/.github/trigger_files/beam_PostCommit_Python_Xlang_IO_Direct.json @@ -1,4 +1,4 @@ { "comment": "Modify this file in a trivial way to cause this test suite to run", - "modification": 1 + "modification": 2 } diff --git a/sdks/java/io/iceberg/build.gradle b/sdks/java/io/iceberg/build.gradle index 0f0fa0a2bb9f..b04ce5d4924c 100644 --- a/sdks/java/io/iceberg/build.gradle +++ b/sdks/java/io/iceberg/build.gradle @@ -41,7 +41,6 @@ hadoopVersions.each {kv -> configurations.create("hadoopVersion$kv.key")} def iceberg_version = "1.9.2" def parquet_version = "1.15.2" -def orc_version = "1.9.2" def hive_version = "3.1.3" dependencies { @@ -53,7 +52,6 @@ dependencies { implementation library.java.joda_time implementation "org.apache.parquet:parquet-column:$parquet_version" implementation "org.apache.parquet:parquet-hadoop:$parquet_version" - implementation "org.apache.orc:orc-core:$orc_version" implementation "org.apache.iceberg:iceberg-core:$iceberg_version" implementation "org.apache.iceberg:iceberg-api:$iceberg_version" implementation "org.apache.iceberg:iceberg-parquet:$iceberg_version" diff --git a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/ScanTaskReader.java b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/ScanTaskReader.java index 81ec229df70f..ac4d1fcdef2d 100644 --- a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/ScanTaskReader.java +++ b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/ScanTaskReader.java @@ -39,7 +39,6 @@ import org.apache.iceberg.data.IdentityPartitionConverters; import org.apache.iceberg.data.Record; import org.apache.iceberg.data.avro.DataReader; -import org.apache.iceberg.data.orc.GenericOrcReader; import org.apache.iceberg.data.parquet.GenericParquetReaders; import org.apache.iceberg.encryption.EncryptionManager; import org.apache.iceberg.encryption.InputFilesDecryptor; @@ -48,7 +47,6 @@ import org.apache.iceberg.io.FileIO; import org.apache.iceberg.io.InputFile; import org.apache.iceberg.mapping.NameMappingParser; -import org.apache.iceberg.orc.ORC; import org.apache.iceberg.parquet.Parquet; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -126,23 +124,6 @@ public boolean advance() throws IOException { CloseableIterable iterable; switch (file.format()) { - case ORC: - LOG.info("Preparing ORC input"); - ORC.ReadBuilder orcReader = - ORC.read(input) - .split(fileTask.start(), fileTask.length()) - .project(requiredSchema) - .createReaderFunc( - fileSchema -> - GenericOrcReader.buildReader(requiredSchema, fileSchema, idToConstants)) - .filter(fileTask.residual()); - - if (nameMapping != null) { - orcReader.withNameMapping(NameMappingParser.fromJson(nameMapping)); - } - - iterable = orcReader.build(); - break; case PARQUET: LOG.info("Preparing Parquet input."); Parquet.ReadBuilder parquetReader = diff --git a/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/TestDataWarehouse.java b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/TestDataWarehouse.java index 61eba3f6ff88..56c629c2f8fe 100644 --- a/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/TestDataWarehouse.java +++ b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/TestDataWarehouse.java @@ -25,6 +25,7 @@ import java.util.Arrays; import java.util.List; import java.util.Map; +import java.util.Objects; import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableMap; import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.Lists; import org.apache.hadoop.conf.Configuration; @@ -42,12 +43,10 @@ import org.apache.iceberg.catalog.Namespace; import org.apache.iceberg.catalog.TableIdentifier; import org.apache.iceberg.data.Record; -import org.apache.iceberg.data.orc.GenericOrcWriter; import org.apache.iceberg.data.parquet.GenericParquetWriter; import org.apache.iceberg.hadoop.HadoopCatalog; import org.apache.iceberg.hadoop.HadoopInputFile; import org.apache.iceberg.io.FileAppender; -import org.apache.iceberg.orc.ORC; import org.apache.iceberg.parquet.Parquet; import org.checkerframework.checker.nullness.qual.Nullable; import org.junit.Assert; @@ -132,25 +131,15 @@ public DataFile writeRecords( FileFormat format = FileFormat.fromFileName(filename); FileAppender appender; - switch (format) { - case PARQUET: - appender = - Parquet.write(fromPath(path, hadoopConf)) - .createWriterFunc(GenericParquetWriter::buildWriter) - .schema(schema) - .overwrite() - .build(); - break; - case ORC: - appender = - ORC.write(fromPath(path, hadoopConf)) - .createWriterFunc(GenericOrcWriter::buildWriter) - .schema(schema) - .overwrite() - .build(); - break; - default: - throw new IOException("Unable to create appender for " + format); + if (Objects.requireNonNull(format) == FileFormat.PARQUET) { + appender = + Parquet.write(fromPath(path, hadoopConf)) + .createWriterFunc(GenericParquetWriter::buildWriter) + .schema(schema) + .overwrite() + .build(); + } else { + throw new IOException("Unable to create appender for " + format); } appender.addAll(records); appender.close();