diff --git a/backends-velox/src-iceberg/test/scala/org/apache/gluten/execution/enhanced/VeloxIcebergSuite.scala b/backends-velox/src-iceberg/test/scala/org/apache/gluten/execution/enhanced/VeloxIcebergSuite.scala index 3a8c48d5e9f..97448cc5e2b 100644 --- a/backends-velox/src-iceberg/test/scala/org/apache/gluten/execution/enhanced/VeloxIcebergSuite.scala +++ b/backends-velox/src-iceberg/test/scala/org/apache/gluten/execution/enhanced/VeloxIcebergSuite.scala @@ -63,6 +63,29 @@ class VeloxIcebergSuite extends IcebergSuite { } } + test("iceberg native ORC write") { + withTable("iceberg_orc_write") { + spark.sql(""" + |create table iceberg_orc_write(a int, b string) using iceberg + |tblproperties ('write.format.default' = 'orc') + |""".stripMargin) + + val df = spark.sql("insert into iceberg_orc_write values (1, 'a'), (2, 'b')") + assert( + df.queryExecution.executedPlan + .asInstanceOf[CommandResultExec] + .commandPhysicalPlan + .isInstanceOf[VeloxIcebergAppendDataExec]) + + checkAnswer( + spark.sql("select * from iceberg_orc_write order by a"), + Seq(Row(1, "a"), Row(2, "b"))) + checkAnswer( + spark.sql("select distinct file_format from default.iceberg_orc_write.files"), + Seq(Row("ORC"))) + } + } + test("iceberg insert partition table identity transform") { withTable("iceberg_tb2") { spark.sql(""" diff --git a/cpp/velox/compute/VeloxBackend.cc b/cpp/velox/compute/VeloxBackend.cc index 0d45559f0e2..da57e4ad9f7 100644 --- a/cpp/velox/compute/VeloxBackend.cc +++ b/cpp/velox/compute/VeloxBackend.cc @@ -40,6 +40,7 @@ #endif #include "compute/VeloxRuntime.h" +#include "compute/iceberg/IcebergWriter.h" #include "config/VeloxConfig.h" #ifdef ENABLE_S3 #include "filesystem/GlutenS3FileSystem.h" @@ -246,6 +247,7 @@ void VeloxBackend::init( velox::parquet::registerParquetReaderFactory(); velox::parquet::registerParquetWriterFactory(); velox::orc::registerOrcReaderFactory(); + registerIcebergOrcWriterFactory(); velox::exec::ExprToSubfieldFilterParser::registerParser(std::make_unique()); velox::connector::hive::BufferedInputBuilder::registerBuilder(std::make_shared()); diff --git a/cpp/velox/compute/iceberg/IcebergWriter.cc b/cpp/velox/compute/iceberg/IcebergWriter.cc index 76d7bb93c48..78b86083aed 100644 --- a/cpp/velox/compute/iceberg/IcebergWriter.cc +++ b/cpp/velox/compute/iceberg/IcebergWriter.cc @@ -26,12 +26,34 @@ #include "utils/ConfigExtractor.h" #include "velox/connectors/hive/iceberg/IcebergDataSink.h" #include "velox/connectors/hive/iceberg/IcebergDeleteFile.h" +#include "velox/dwio/common/WriterFactory.h" +#include "velox/dwio/dwrf/writer/Writer.h" using namespace facebook::velox; using namespace facebook::velox::connector::hive; using namespace facebook::velox::connector::hive::iceberg; namespace { +class IcebergOrcWriterFactory final : public dwio::common::WriterFactory { + public: + IcebergOrcWriterFactory() : WriterFactory(dwio::common::FileFormat::ORC) {} + + std::unique_ptr createWriter( + std::unique_ptr sink, + const std::shared_ptr& options) override { + auto orcOptions = std::dynamic_pointer_cast(options); + VELOX_CHECK_NOT_NULL(orcOptions, "Iceberg ORC writer expected DWRF writer options."); + VELOX_CHECK_EQ(orcOptions->format, dwrf::DwrfFormat::kOrc); + return std::make_unique(std::move(sink), *orcOptions); + } + + std::unique_ptr createWriterOptions() override { + auto options = std::make_unique(); + options->format = dwrf::DwrfFormat::kOrc; + return options; + } +}; + // Custom Iceberg file name generator for Gluten class GlutenIcebergFileNameGenerator : public connector::hive::FileNameGenerator { public: @@ -175,6 +197,10 @@ std::shared_ptr createIcebergInsertTableHandle( } // namespace namespace gluten { +void registerIcebergOrcWriterFactory() { + dwio::common::registerWriterFactory(std::make_shared()); +} + IcebergWriter::IcebergWriter( const RowTypePtr& rowType, int32_t format, diff --git a/cpp/velox/compute/iceberg/IcebergWriter.h b/cpp/velox/compute/iceberg/IcebergWriter.h index 0ab3a803608..641cb9450eb 100644 --- a/cpp/velox/compute/iceberg/IcebergWriter.h +++ b/cpp/velox/compute/iceberg/IcebergWriter.h @@ -25,6 +25,8 @@ namespace gluten { +void registerIcebergOrcWriterFactory(); + struct WriteStats { uint64_t numWrittenBytes{0}; uint32_t numWrittenFiles{0}; diff --git a/docs/get-started/VeloxIceberg.md b/docs/get-started/VeloxIceberg.md index aea2e89c76f..1ec7a6f9cd4 100644 --- a/docs/get-started/VeloxIceberg.md +++ b/docs/get-started/VeloxIceberg.md @@ -160,7 +160,7 @@ The "Gluten Support" column is now ready to be populated with: | Spark option | Default | Description | Gluten Support | | --- | --- | --- | --- | -| write-format | Table write.format.default | File format to use for this write operation; parquet, avro, or orc |⚠️ Parquet only| +| write-format | Table write.format.default | File format to use for this write operation; parquet, avro, or orc |⚠️ Parquet and ORC| | target-file-size-bytes | As per table property | Overrides this table's write.target-file-size-bytes | | | check-nullability | true | Sets the nullable check on fields | | | snapshot-property.custom-key | null | Adds an entry with custom-key and corresponding value in the snapshot summary (the snapshot-property. prefix is only required for DSv2) | | @@ -194,7 +194,7 @@ extracted from https://iceberg.apache.org/docs/latest/configuration/ | Property | Default | Description | Gluten Support | | --- | --- | --- | --- | -| write.format.default | parquet | Default file format for the table; parquet, avro, or orc | | +| write.format.default | parquet | Default file format for the table; parquet, avro, or orc |⚠️ Parquet and ORC| | write.delete.format.default | data file format | Default delete file format for the table; parquet, avro, or orc | | | write.parquet.row-group-size-bytes | 134217728 (128 MB) | Parquet row group size | | | write.parquet.page-size-bytes | 1048576 (1 MB) | Parquet page size |✅| diff --git a/gluten-iceberg/src/main/scala/org/apache/gluten/execution/IcebergWriteExec.scala b/gluten-iceberg/src/main/scala/org/apache/gluten/execution/IcebergWriteExec.scala index 270745063af..b4db719386d 100644 --- a/gluten-iceberg/src/main/scala/org/apache/gluten/execution/IcebergWriteExec.scala +++ b/gluten-iceberg/src/main/scala/org/apache/gluten/execution/IcebergWriteExec.scala @@ -108,7 +108,7 @@ trait IcebergWriteExec extends ColumnarV2TableWriteExec { return ValidationResult.failed("Not support write table with sort order") } val format = IcebergWriteUtil.getFileFormat(write) - if (format != FileFormat.PARQUET) { + if (!Seq(FileFormat.PARQUET, FileFormat.ORC).contains(format)) { return ValidationResult.failed("Not support this format " + format.name()) }