Skip to content

Commit ddb87cc

Browse files
committed
✨ feat(data): Apply write defaults in writers
DataWriter aligns batches to the write schema when the input schema differs, via a new Write(data, input_schema) overload; schemas that already match keep the zero-conversion Write(data) path. Adds Parquet and Avro end-to-end tests where a column added with write-default reads back the default for rows written without it.
1 parent 3486a4c commit ddb87cc

9 files changed

Lines changed: 156 additions & 2 deletions

File tree

‎src/iceberg/avro/avro_writer.cc‎

Lines changed: 13 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_c_data_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"
@@ -241,6 +242,14 @@ class AvroWriter::Impl {
241242
return {};
242243
}
243244

245+
Status Write(ArrowArray* data, const Schema& input_schema) {
246+
ICEBERG_ASSIGN_OR_RAISE(
247+
auto aligned_data,
248+
arrow::AlignBatchForWrite(data, input_schema, *write_schema_,
249+
::arrow::default_memory_pool()));
250+
return Write(&aligned_data);
251+
}
252+
244253
Status Close() {
245254
if (!backend_->Closed()) {
246255
backend_->Close();
@@ -290,6 +299,10 @@ AvroWriter::~AvroWriter() = default;
290299

291300
Status AvroWriter::Write(ArrowArray* data) { return impl_->Write(data); }
292301

302+
Status AvroWriter::Write(ArrowArray* data, const Schema& input_schema) {
303+
return impl_->Write(data, input_schema);
304+
}
305+
293306
Status AvroWriter::Open(const WriterOptions& options) {
294307
impl_ = std::make_unique<Impl>();
295308
return impl_->Open(options);

‎src/iceberg/avro/avro_writer.h‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -40,6 +40,8 @@ class ICEBERG_BUNDLE_EXPORT AvroWriter : public Writer {
4040

4141
Status Write(ArrowArray* data) final;
4242

43+
Status Write(ArrowArray* data, const Schema& input_schema) final;
44+
4345
Result<Metrics> metrics() final;
4446

4547
Result<int64_t> length() final;

‎src/iceberg/data/data_writer.cc‎

Lines changed: 9 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -26,13 +26,16 @@
2626
#include "iceberg/file_writer.h"
2727
#include "iceberg/manifest/manifest_entry.h"
2828
#include "iceberg/partition_spec.h"
29+
#include "iceberg/schema.h"
2930
#include "iceberg/util/macros.h"
3031

3132
namespace iceberg {
3233

3334
class DataWriter::Impl {
3435
public:
3536
static Result<std::unique_ptr<Impl>> Make(DataWriterOptions options) {
37+
ICEBERG_PRECHECK(options.schema != nullptr, "Write schema must not be null");
38+
ICEBERG_PRECHECK(options.input_schema != nullptr, "Input schema must not be null");
3639
WriterOptions writer_options{
3740
.path = options.path,
3841
.schema = options.schema,
@@ -46,7 +49,12 @@ class DataWriter::Impl {
4649
return std::unique_ptr<Impl>(new Impl(std::move(options), std::move(writer)));
4750
}
4851

49-
Status Write(ArrowArray* data) { return writer_->Write(data); }
52+
Status Write(ArrowArray* data) {
53+
if (*options_.input_schema == *options_.schema) {
54+
return writer_->Write(data);
55+
}
56+
return writer_->Write(data, *options_.input_schema);
57+
}
5058

5159
Result<int64_t> Length() const { return writer_->length(); }
5260

‎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+
/// Required schema of Arrow batches passed to Write.
46+
std::shared_ptr<Schema> input_schema;
4547
std::shared_ptr<PartitionSpec> spec;
4648
PartitionValues partition;
4749
FileFormatType format = FileFormatType::kParquet;

‎src/iceberg/file_writer.h‎

Lines changed: 7 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -75,7 +75,7 @@ class ICEBERG_EXPORT WriterProperties : public ConfigBase<WriterProperties> {
7575
struct ICEBERG_EXPORT WriterOptions {
7676
/// \brief The path to the file to write.
7777
std::string path;
78-
/// \brief The schema of the data to write.
78+
/// \brief The target schema to write to the file.
7979
std::shared_ptr<Schema> schema;
8080
/// \brief FileIO instance to create the file.
8181
std::shared_ptr<class FileIO> io;
@@ -107,6 +107,12 @@ class ICEBERG_EXPORT Writer {
107107
/// \note Ownership of the data is transferred to the writer.
108108
virtual Status Write(ArrowArray* data) = 0;
109109

110+
/// \brief Write Arrow data described by an input schema that may differ from the
111+
/// writer schema.
112+
///
113+
/// \note Ownership of the data is transferred to the writer.
114+
virtual Status Write(ArrowArray* data, const Schema& input_schema) = 0;
115+
110116
/// \brief Get the file statistics.
111117
/// Only valid after the file is closed.
112118
virtual Result<Metrics> metrics() = 0;

‎src/iceberg/parquet/parquet_writer.cc‎

Lines changed: 12 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -39,6 +39,7 @@
3939

4040
#include "iceberg/arrow/arrow_io_internal.h"
4141
#include "iceberg/arrow/arrow_status_internal.h"
42+
#include "iceberg/arrow_c_data_util_internal.h"
4243
#include "iceberg/parquet/parquet_metrics_internal.h"
4344
#include "iceberg/parquet/parquet_schema_util_internal.h"
4445
#include "iceberg/schema_internal.h"
@@ -295,6 +296,13 @@ class ParquetWriter::Impl {
295296
return {};
296297
}
297298

299+
Status Write(ArrowArray* array, const Schema& input_schema) {
300+
ICEBERG_ASSIGN_OR_RAISE(
301+
auto aligned_array,
302+
arrow::AlignBatchForWrite(array, input_schema, *schema_, pool_));
303+
return Write(&aligned_array);
304+
}
305+
298306
// Close the writer and release resources
299307
Status Close() {
300308
if (writer_ == nullptr) {
@@ -371,6 +379,10 @@ Status ParquetWriter::Open(const WriterOptions& options) {
371379

372380
Status ParquetWriter::Write(ArrowArray* array) { return impl_->Write(array); }
373381

382+
Status ParquetWriter::Write(ArrowArray* array, const Schema& input_schema) {
383+
return impl_->Write(array, input_schema);
384+
}
385+
374386
Status ParquetWriter::Close() { return impl_->Close(); }
375387

376388
Result<Metrics> ParquetWriter::metrics() { return impl_->metrics(); }

‎src/iceberg/parquet/parquet_writer.h‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -40,6 +40,8 @@ class ICEBERG_BUNDLE_EXPORT ParquetWriter : public Writer {
4040

4141
Status Write(ArrowArray* array) final;
4242

43+
Status Write(ArrowArray* array, const Schema& input_schema) final;
44+
4345
Result<Metrics> metrics() final;
4446

4547
Result<int64_t> length() final;

‎src/iceberg/test/data_writer_test.cc‎

Lines changed: 13 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -71,6 +71,7 @@ class DataWriterTest : public ::testing::Test {
7171
return DataWriterOptions{
7272
.path = "test_data.parquet",
7373
.schema = schema_,
74+
.input_schema = schema_,
7475
.spec = partition_spec_,
7576
.partition = std::move(partition),
7677
.format = FileFormatType::kParquet,
@@ -138,11 +139,22 @@ class DataWriterFormatTest
138139
: public DataWriterTest,
139140
public ::testing::WithParamInterface<std::pair<FileFormatType, std::string>> {};
140141

142+
TEST_F(DataWriterTest, RejectsMissingInputSchema) {
143+
auto options = MakeDefaultOptions();
144+
options.input_schema.reset();
145+
146+
auto result = DataWriter::Make(options);
147+
148+
EXPECT_THAT(result, IsError(ErrorKind::kInvalidArgument));
149+
EXPECT_THAT(result, HasErrorMessage("Input schema must not be null"));
150+
}
151+
141152
TEST_P(DataWriterFormatTest, CreateWithFormat) {
142153
auto [format, path] = GetParam();
143154
DataWriterOptions options{
144155
.path = path,
145156
.schema = schema_,
157+
.input_schema = schema_,
146158
.spec = partition_spec_,
147159
.partition = PartitionValues{},
148160
.format = format,
@@ -162,6 +174,7 @@ TEST_P(DataWriterFormatTest, WriteRowLineage) {
162174
DataWriterOptions options{
163175
.path = path,
164176
.schema = schema,
177+
.input_schema = schema,
165178
.spec = partition_spec_,
166179
.partition = PartitionValues{},
167180
.format = format,

‎src/iceberg/test/default_value_test.cc‎

Lines changed: 96 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -23,6 +23,8 @@
2323
#include <memory>
2424
#include <string>
2525
#include <string_view>
26+
#include <unordered_map>
27+
#include <utility>
2628
#include <vector>
2729

2830
#include <arrow/array.h>
@@ -38,11 +40,14 @@
3840

3941
#include "iceberg/arrow/arrow_io_internal.h"
4042
#include "iceberg/avro/avro_register.h"
43+
#include "iceberg/data/data_writer.h"
4144
#include "iceberg/expression/literal.h"
4245
#include "iceberg/file_format.h"
4346
#include "iceberg/file_reader.h"
4447
#include "iceberg/file_writer.h"
4548
#include "iceberg/parquet/parquet_register.h"
49+
#include "iceberg/partition_spec.h"
50+
#include "iceberg/row/partition_values.h"
4651
#include "iceberg/schema.h"
4752
#include "iceberg/schema_field.h"
4853
#include "iceberg/schema_internal.h"
@@ -72,6 +77,12 @@ struct DefaultValueEndToEndParam {
7277
bool avro_skip_datum = true;
7378
};
7479

80+
struct WriteDefaultEndToEndParam {
81+
std::string name;
82+
FileFormatType format;
83+
std::string path;
84+
};
85+
7586
class DefaultValueEndToEndTest
7687
: public UpdateTestBase,
7788
public ::testing::WithParamInterface<DefaultValueEndToEndParam> {
@@ -86,6 +97,37 @@ class DefaultValueEndToEndTest
8697
}
8798
};
8899

100+
class WriteDefaultEndToEndTest
101+
: public ::testing::TestWithParam<WriteDefaultEndToEndParam> {
102+
protected:
103+
static void SetUpTestSuite() {
104+
parquet::RegisterAll();
105+
avro::RegisterAll();
106+
}
107+
108+
void SetUp() override { file_io_ = arrow::ArrowFileSystemFileIO::MakeMockFileIO(); }
109+
110+
std::shared_ptr<::arrow::Array> CreateArray(const Schema& schema,
111+
std::string_view json) {
112+
ArrowSchema arrow_c_schema;
113+
ICEBERG_THROW_NOT_OK(ToArrowSchema(schema, &arrow_c_schema));
114+
auto arrow_schema = ::arrow::ImportType(&arrow_c_schema).ValueOrDie();
115+
return ::arrow::json::ArrayFromJSONString(::arrow::struct_(arrow_schema->fields()),
116+
std::string(json))
117+
.ValueOrDie();
118+
}
119+
120+
static std::unordered_map<std::string, std::string> FormatProperties(
121+
FileFormatType format) {
122+
if (format == FileFormatType::kParquet) {
123+
return {{"write.parquet.compression-codec", "uncompressed"}};
124+
}
125+
return {};
126+
}
127+
128+
std::shared_ptr<FileIO> file_io_;
129+
};
130+
89131
TEST_P(DefaultValueEndToEndTest, WriteEvolveReadFillsInitialDefault) {
90132
const auto& param = GetParam();
91133
ICEBERG_UNWRAP_OR_FAIL(auto original_schema, table_->schema());
@@ -186,6 +228,50 @@ TEST_P(DefaultValueEndToEndTest, WriteEvolveReadFillsInitialDefault) {
186228
ASSERT_FALSE(next_batch.has_value());
187229
}
188230

231+
TEST_P(WriteDefaultEndToEndTest, MissingColumnUsesWriteDefault) {
232+
const auto& param = GetParam();
233+
auto input_schema = std::make_shared<Schema>(
234+
std::vector<SchemaField>{SchemaField::MakeRequired(1, "id", int32())});
235+
auto write_schema = std::make_shared<Schema>(std::vector<SchemaField>{
236+
SchemaField::MakeRequired(1, "id", int32()),
237+
SchemaField(2, "added", int32(), /*optional=*/false, /*doc=*/{},
238+
std::make_shared<const Literal>(Literal::Int(42)),
239+
std::make_shared<const Literal>(Literal::Int(7))),
240+
});
241+
DataWriterOptions options{
242+
.path = param.path,
243+
.schema = write_schema,
244+
.input_schema = input_schema,
245+
.spec = PartitionSpec::Unpartitioned(),
246+
.partition = PartitionValues{},
247+
.format = param.format,
248+
.io = file_io_,
249+
.properties = FormatProperties(param.format),
250+
};
251+
252+
ICEBERG_UNWRAP_OR_FAIL(auto writer, DataWriter::Make(options));
253+
auto input = CreateArray(*input_schema, R"([[1], [2]])");
254+
ArrowArray arrow_array;
255+
ASSERT_TRUE(::arrow::ExportArray(*input, &arrow_array).ok());
256+
ASSERT_THAT(writer->Write(&arrow_array), IsOk());
257+
ASSERT_THAT(writer->Close(), IsOk());
258+
259+
ICEBERG_UNWRAP_OR_FAIL(
260+
auto reader, ReaderFactoryRegistry::Open(
261+
param.format,
262+
{.path = param.path, .io = file_io_, .projection = write_schema}));
263+
ICEBERG_UNWRAP_OR_FAIL(auto batch, reader->Next());
264+
ASSERT_TRUE(batch.has_value());
265+
266+
auto expected = CreateArray(*write_schema, R"([[1, 7], [2, 7]])");
267+
auto actual = ::arrow::ImportArray(&batch.value(), expected->type()).ValueOrDie();
268+
269+
// This file is written after the column exists, so the missing input column uses
270+
// write-default (7), not initial-default (42).
271+
ASSERT_TRUE(actual->Equals(expected))
272+
<< "actual: " << actual->ToString() << "\nexpected: " << expected->ToString();
273+
}
274+
189275
namespace {
190276

191277
// Two-row column of `json` values at `type`, for the simple cases.
@@ -261,6 +347,16 @@ INSTANTIATE_TEST_SUITE_P(
261347
return info.param.name;
262348
});
263349

350+
INSTANTIATE_TEST_SUITE_P(
351+
Formats, WriteDefaultEndToEndTest,
352+
::testing::Values(WriteDefaultEndToEndParam{"parquet", FileFormatType::kParquet,
353+
"write-default.parquet"},
354+
WriteDefaultEndToEndParam{"avro", FileFormatType::kAvro,
355+
"write-default.avro"}),
356+
[](const ::testing::TestParamInfo<WriteDefaultEndToEndParam>& info) {
357+
return info.param.name;
358+
});
359+
264360
} // namespace
265361

266362
} // namespace iceberg

0 commit comments

Comments
 (0)