Skip to content

Commit 201e6be

Browse files
committed
✨ feat(data): Apply defaults in data writers
Pass the input schema through DataWriter so Parquet and Avro writers can align missing columns with write defaults. Add an end-to-end test that distinguishes write-default from initial-default for both formats. Refs #637
1 parent 6184690 commit 201e6be

5 files changed

Lines changed: 63 additions & 0 deletions

File tree

‎src/iceberg/avro/avro_writer.cc‎

Lines changed: 11 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -32,6 +32,7 @@
3232

3333
#include "iceberg/arrow/arrow_io_internal.h"
3434
#include "iceberg/arrow/arrow_status_internal.h"
35+
#include "iceberg/arrow/write_schema_util_internal.h"
3536
#include "iceberg/avro/avro_data_util_internal.h"
3637
#include "iceberg/avro/avro_direct_encoder_internal.h"
3738
#include "iceberg/avro/avro_metrics.h"
@@ -181,6 +182,7 @@ class AvroWriter::Impl {
181182
public:
182183
Status Open(const WriterOptions& options) {
183184
write_schema_ = options.schema;
185+
input_schema_ = options.input_schema;
184186

185187
::avro::NodePtr root;
186188
ICEBERG_RETURN_UNEXPECTED(ToAvroNodeVisitor{}.Visit(*write_schema_, &root));
@@ -229,6 +231,13 @@ class AvroWriter::Impl {
229231
}
230232

231233
Status Write(ArrowArray* data) {
234+
ArrowArray aligned_data{};
235+
if (input_schema_ != nullptr) {
236+
ICEBERG_ASSIGN_OR_RAISE(
237+
aligned_data, arrow::AlignBatchForWrite(data, *input_schema_, *write_schema_,
238+
::arrow::default_memory_pool()));
239+
data = &aligned_data;
240+
}
232241
ICEBERG_ARROW_ASSIGN_OR_RETURN(auto batch,
233242
::arrow::ImportRecordBatch(data, arrow_schema_));
234243

@@ -272,6 +281,8 @@ class AvroWriter::Impl {
272281
private:
273282
// The schema to write.
274283
std::shared_ptr<::iceberg::Schema> write_schema_;
284+
// Schema of incoming Arrow batches. Null means batches already match write_schema_.
285+
std::shared_ptr<::iceberg::Schema> input_schema_;
275286
// The avro schema to write.
276287
std::shared_ptr<::avro::ValidSchema> avro_schema_;
277288
// Arrow output stream of the Avro file to write

‎src/iceberg/data/data_writer.cc‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -36,6 +36,7 @@ class DataWriter::Impl {
3636
WriterOptions writer_options{
3737
.path = options.path,
3838
.schema = options.schema,
39+
.input_schema = options.input_schema,
3940
.io = options.io,
4041
.properties = WriterProperties::FromMap(options.properties),
4142
};

‎src/iceberg/data/data_writer.h‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -42,6 +42,8 @@ namespace iceberg {
4242
struct ICEBERG_DATA_EXPORT DataWriterOptions {
4343
std::string path;
4444
std::shared_ptr<Schema> schema;
45+
/// Schema of Arrow batches passed to Write. When null, batches must match `schema`.
46+
std::shared_ptr<Schema> input_schema;
4547
std::shared_ptr<PartitionSpec> spec;
4648
PartitionValues partition;
4749
FileFormatType format = FileFormatType::kParquet;

‎src/iceberg/parquet/parquet_writer.cc‎

Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -40,6 +40,7 @@
4040

4141
#include "iceberg/arrow/arrow_io_internal.h"
4242
#include "iceberg/arrow/arrow_status_internal.h"
43+
#include "iceberg/arrow/write_schema_util_internal.h"
4344
#include "iceberg/parquet/parquet_metrics_internal.h"
4445
#include "iceberg/schema_internal.h"
4546
#include "iceberg/type.h"
@@ -239,6 +240,7 @@ class ParquetWriter::Impl {
239240
public:
240241
Status Open(const WriterOptions& options) {
241242
schema_ = options.schema;
243+
input_schema_ = options.input_schema;
242244

243245
ICEBERG_ASSIGN_OR_RAISE(auto compression, ParseCompression(options.properties));
244246
ICEBERG_ASSIGN_OR_RAISE(auto compression_level, ParseCodecLevel(options.properties));
@@ -283,6 +285,12 @@ class ParquetWriter::Impl {
283285
}
284286

285287
Status Write(ArrowArray* array) {
288+
ArrowArray aligned_array{};
289+
if (input_schema_ != nullptr) {
290+
ICEBERG_ASSIGN_OR_RAISE(aligned_array, arrow::AlignBatchForWrite(
291+
array, *input_schema_, *schema_, pool_));
292+
array = &aligned_array;
293+
}
286294
ICEBERG_ARROW_ASSIGN_OR_RETURN(auto batch,
287295
::arrow::ImportRecordBatch(array, arrow_schema_));
288296

@@ -344,6 +352,8 @@ class ParquetWriter::Impl {
344352
::arrow::MemoryPool* pool_ = ::arrow::default_memory_pool();
345353
// Schema to write from the Iceberg table.
346354
std::shared_ptr<Schema> schema_;
355+
// Schema of incoming Arrow batches. Null means batches already match schema_.
356+
std::shared_ptr<Schema> input_schema_;
347357
// Schema to write from the Parquet file.
348358
std::shared_ptr<::arrow::Schema> arrow_schema_;
349359
// Parquet schema descriptor generated from the Arrow schema.

‎src/iceberg/test/data_writer_test.cc‎

Lines changed: 39 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -29,6 +29,7 @@
2929
#include "iceberg/avro/avro_register.h"
3030
#include "iceberg/data/equality_delete_writer.h"
3131
#include "iceberg/data/position_delete_writer.h"
32+
#include "iceberg/expression/literal.h"
3233
#include "iceberg/file_format.h"
3334
#include "iceberg/file_reader.h"
3435
#include "iceberg/manifest/manifest_entry.h"
@@ -188,6 +189,44 @@ TEST_P(DataWriterFormatTest, WriteRowLineage) {
188189
[3, 102, 7]])"));
189190
}
190191

192+
TEST_P(DataWriterFormatTest, WriteMissingColumnUsesWriteDefault) {
193+
auto [format, path] = GetParam();
194+
auto input_schema = std::make_shared<Schema>(
195+
std::vector<SchemaField>{SchemaField::MakeRequired(1, "id", int32())});
196+
auto write_schema = std::make_shared<Schema>(std::vector<SchemaField>{
197+
SchemaField::MakeRequired(1, "id", int32()),
198+
SchemaField(2, "added", int32(), /*optional=*/false, /*doc=*/{},
199+
std::make_shared<const Literal>(Literal::Int(42)),
200+
std::make_shared<const Literal>(Literal::Int(7))),
201+
});
202+
DataWriterOptions options{
203+
.path = path,
204+
.schema = write_schema,
205+
.input_schema = input_schema,
206+
.spec = partition_spec_,
207+
.partition = PartitionValues{},
208+
.format = format,
209+
.io = file_io_,
210+
.properties = FormatProperties(format),
211+
};
212+
213+
ICEBERG_UNWRAP_OR_FAIL(auto writer, DataWriter::Make(options));
214+
auto input = CreateArray(*input_schema, R"([[1], [2]])");
215+
ArrowArray arrow_array;
216+
ASSERT_TRUE(::arrow::ExportArray(*input, &arrow_array).ok());
217+
ASSERT_THAT(writer->Write(&arrow_array), IsOk());
218+
ASSERT_THAT(writer->Close(), IsOk());
219+
220+
ICEBERG_UNWRAP_OR_FAIL(
221+
auto reader,
222+
ReaderFactoryRegistry::Open(
223+
format, {.path = path, .io = file_io_, .projection = write_schema}));
224+
225+
// The file was written after the column was added, so the missing input column uses
226+
// write-default (7), not initial-default (42).
227+
ASSERT_NO_FATAL_FAILURE(VerifyNextBatch(*reader, *write_schema, R"([[1, 7], [2, 7]])"));
228+
}
229+
191230
INSTANTIATE_TEST_SUITE_P(
192231
FormatTypes, DataWriterFormatTest,
193232
::testing::Values(std::make_pair(FileFormatType::kParquet, "test_data.parquet"),

0 commit comments

Comments
 (0)