Skip to content

Commit 3b56fb3

Browse files
committed
fix(blob): support blob-only writes and reject nested blob fields
1 parent 6c88298 commit 3b56fb3

11 files changed

Lines changed: 239 additions & 59 deletions

src/paimon/common/data/blob_utils.cpp

Lines changed: 4 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -82,13 +82,12 @@ Result<BlobUtils::SeparatedStructArrays> BlobUtils::SeparateBlobArray(
8282
return Status::Invalid(
8383
"SeparateBlobArray expects at least one non-inline blob field, but got none.");
8484
}
85-
if (main_fields.empty()) {
86-
return Status::Invalid("SeparateBlobArray expects at least one main field, but got none.");
87-
}
8885

8986
SeparatedStructArrays result;
90-
PAIMON_ASSIGN_OR_RAISE_FROM_ARROW(result.main_array,
91-
arrow::StructArray::Make(main_arrays, main_fields));
87+
if (!main_fields.empty()) {
88+
PAIMON_ASSIGN_OR_RAISE_FROM_ARROW(result.main_array,
89+
arrow::StructArray::Make(main_arrays, main_fields));
90+
}
9291
PAIMON_ASSIGN_OR_RAISE_FROM_ARROW(result.blob_array,
9392
arrow::StructArray::Make(blob_arrays, blob_fields));
9493
return result;

src/paimon/common/data/blob_utils.h

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -51,7 +51,8 @@ class PAIMON_EXPORT BlobUtils {
5151
};
5252

5353
struct SeparatedStructArrays {
54-
/// Non-blob fields (includes inline blob fields when inline_fields is provided)
54+
/// Non-blob fields (includes inline blob fields when inline_fields is provided).
55+
/// nullptr when all fields are stored in blob files.
5556
std::shared_ptr<arrow::StructArray> main_array;
5657
/// Blob fields that go to separate .blob files
5758
std::shared_ptr<arrow::StructArray> blob_array;

src/paimon/common/data/blob_utils_test.cpp

Lines changed: 5 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -182,11 +182,13 @@ TEST_F(BlobUtilsTest, SeparateBlobArray) {
182182
BlobUtils::SeparateBlobArray(struct_array, /*inline_fields=*/{"f2_blob"}),
183183
"SeparateBlobArray expects at least one non-inline blob field, but got none.");
184184

185-
// All fields are blob with no inline -> no main field -> should return error
185+
// All fields are blob with no inline -> no main array is needed
186186
auto all_blob_struct = arrow::StructArray::Make({blob_array_data}, {blob_field}).ValueOrDie();
187187
auto all_blob_sa = std::dynamic_pointer_cast<arrow::StructArray>(all_blob_struct);
188-
ASSERT_NOK_WITH_MSG(BlobUtils::SeparateBlobArray(all_blob_sa, /*inline_fields=*/{}),
189-
"SeparateBlobArray expects at least one main field, but got none.");
188+
ASSERT_OK_AND_ASSIGN(auto all_blob_separated,
189+
BlobUtils::SeparateBlobArray(all_blob_sa, /*inline_fields=*/{}));
190+
ASSERT_EQ(nullptr, all_blob_separated.main_array);
191+
ASSERT_TRUE(all_blob_separated.blob_array->Equals(*all_blob_sa));
190192
}
191193

192194
TEST_F(BlobUtilsTest, SeparateBlobArrayWithPartialInline) {

src/paimon/core/append/append_only_writer.cpp

Lines changed: 7 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -238,10 +238,14 @@ AppendOnlyWriter::RollingFileWriterResult AppendOnlyWriter::CreateRollingBlobWri
238238
options_.GetBlobTargetFileSize(), single_blob_file_writer_factory);
239239
};
240240

241+
WriterFactory main_writer_factory;
242+
if (schemas.main_schema->num_fields() > 0) {
243+
main_writer_factory =
244+
GetDataFileWriterFactory(schemas.main_schema, schemas.main_schema->field_names());
245+
}
241246
return std::make_unique<RollingBlobFileWriter>(
242-
options_.GetTargetFileSize(/*has_primary_key=*/false),
243-
GetDataFileWriterFactory(schemas.main_schema, schemas.main_schema->field_names()),
244-
blob_schema, blob_writer_creator, arrow::struct_(write_schema_->fields()), inline_fields);
247+
options_.GetTargetFileSize(/*has_primary_key=*/false), main_writer_factory, blob_schema,
248+
blob_writer_creator, arrow::struct_(write_schema_->fields()), inline_fields);
245249
}
246250

247251
Status AppendOnlyWriter::Sync() {

src/paimon/core/append/append_only_writer_test.cpp

Lines changed: 34 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -741,6 +741,40 @@ TEST_F(AppendOnlyWriterTest, TestWriteWithSingleBlobField) {
741741
ASSERT_OK(writer->Close());
742742
}
743743

744+
TEST_F(AppendOnlyWriterTest, TestWriteWithOnlyBlobField) {
745+
auto options =
746+
CreateOptions({{Options::FILE_FORMAT, "orc"}, {Options::MANIFEST_FORMAT, "orc"}});
747+
auto dir = UniqueTestDirectory::Create();
748+
ASSERT_TRUE(dir);
749+
auto path_factory = CreatePathFactory(dir->Str(), "orc", options);
750+
751+
auto blob_field = BlobUtils::ToArrowField("blob", false);
752+
auto schema = arrow::schema({blob_field});
753+
ASSERT_OK_AND_ASSIGN(auto writer,
754+
CreateAppendOnlyWriter(options, /*schema_id=*/0, schema,
755+
/*write_cols=*/std::vector<std::string>{"blob"},
756+
/*max_sequence_number=*/-1, path_factory,
757+
compact_manager_, memory_pool_));
758+
759+
arrow::LargeBinaryBuilder blob_builder;
760+
ASSERT_TRUE(blob_builder.Append("a", 1).ok());
761+
ASSERT_TRUE(blob_builder.Append("bb", 2).ok());
762+
auto blob_array = blob_builder.Finish().ValueOrDie();
763+
764+
ASSERT_OK(writer->Write(CreateStructBatch(schema, {blob_array})));
765+
ASSERT_OK_AND_ASSIGN(CommitIncrement inc, writer->PrepareCommit(/*wait_compaction=*/true));
766+
767+
const auto& new_files = inc.GetNewFilesIncrement().NewFiles();
768+
ASSERT_EQ(new_files.size(), 1);
769+
ASSERT_TRUE(BlobUtils::IsBlobFile(new_files[0]->file_name));
770+
ASSERT_EQ(new_files[0]->row_count, 2);
771+
ASSERT_TRUE(new_files[0]->write_cols.has_value());
772+
ASSERT_EQ(new_files[0]->write_cols.value(), std::vector<std::string>({"blob"}));
773+
ASSERT_TRUE(
774+
options.GetFileSystem()->Exists(path_factory->ToPath(new_files[0]->file_name)).value());
775+
ASSERT_OK(writer->Close());
776+
}
777+
744778
TEST_F(AppendOnlyWriterTest, TestWriteWithMultipleBlobFields) {
745779
auto options =
746780
CreateOptions({{Options::FILE_FORMAT, "orc"}, {Options::MANIFEST_FORMAT, "orc"}});

src/paimon/core/io/rolling_blob_file_writer.cpp

Lines changed: 24 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -55,7 +55,7 @@ RollingBlobFileWriter::RollingBlobFileWriter(
5555
Status RollingBlobFileWriter::Write(::ArrowArray* record) {
5656
ScopeGuard guard([this]() -> void { this->Abort(); });
5757
// Open the current writer if write the first record or roll over happen before.
58-
if (PAIMON_UNLIKELY(current_writer_ == nullptr)) {
58+
if (writer_factory_ != nullptr && PAIMON_UNLIKELY(current_writer_ == nullptr)) {
5959
PAIMON_RETURN_NOT_OK(OpenCurrentWriter());
6060
}
6161
if (PAIMON_UNLIKELY(blob_writer_ == nullptr)) {
@@ -69,12 +69,14 @@ Status RollingBlobFileWriter::Write(::ArrowArray* record) {
6969
PAIMON_ASSIGN_OR_RAISE(BlobUtils::SeparatedStructArrays separated_arrays,
7070
BlobUtils::SeparateBlobArray(struct_array, inline_fields_));
7171
// Write main (non-blob) data
72-
::ArrowArray c_main_array;
73-
PAIMON_RETURN_NOT_OK_FROM_ARROW(
74-
arrow::ExportArray(*separated_arrays.main_array, &c_main_array));
75-
ScopeGuard array_lifecycle_guard(
76-
[&c_main_array]() -> void { ArrowArrayRelease(&c_main_array); });
77-
PAIMON_RETURN_NOT_OK(current_writer_->Write(&c_main_array));
72+
if (current_writer_ != nullptr) {
73+
::ArrowArray c_main_array;
74+
PAIMON_RETURN_NOT_OK_FROM_ARROW(
75+
arrow::ExportArray(*separated_arrays.main_array, &c_main_array));
76+
ScopeGuard array_lifecycle_guard(
77+
[&c_main_array]() -> void { ArrowArrayRelease(&c_main_array); });
78+
PAIMON_RETURN_NOT_OK(current_writer_->Write(&c_main_array));
79+
}
7880

7981
// Write blob data via MultipleBlobFileWriter (each blob field independently)
8082
::ArrowArray c_blob_array;
@@ -84,28 +86,35 @@ Status RollingBlobFileWriter::Write(::ArrowArray* record) {
8486
PAIMON_RETURN_NOT_OK(blob_writer_->Write(&c_blob_array));
8587

8688
record_count_ += record_count;
87-
PAIMON_ASSIGN_OR_RAISE(bool need_rolling_file, NeedRollingFile());
88-
if (need_rolling_file) {
89-
PAIMON_RETURN_NOT_OK(CloseCurrentWriter());
89+
if (current_writer_ != nullptr) {
90+
PAIMON_ASSIGN_OR_RAISE(bool need_rolling_file, NeedRollingFile());
91+
if (need_rolling_file) {
92+
PAIMON_RETURN_NOT_OK(CloseCurrentWriter());
93+
}
9094
}
9195
guard.Release();
9296
return Status::OK();
9397
}
9498

9599
Status RollingBlobFileWriter::CloseCurrentWriter() {
96-
if (current_writer_ == nullptr) {
100+
if (current_writer_ == nullptr && blob_writer_ == nullptr) {
97101
return Status::OK();
98102
}
99103
if (blob_writer_ == nullptr) {
100104
return Status::OK();
101105
}
102-
PAIMON_ASSIGN_OR_RAISE(std::shared_ptr<DataFileMeta> main_data_file_meta, CloseMainWriter());
106+
std::shared_ptr<DataFileMeta> main_data_file_meta;
107+
if (current_writer_ != nullptr) {
108+
PAIMON_ASSIGN_OR_RAISE(main_data_file_meta, CloseMainWriter());
109+
}
103110
PAIMON_ASSIGN_OR_RAISE(std::vector<std::shared_ptr<DataFileMeta>> blob_metas,
104111
CloseBlobWriter());
105-
PAIMON_RETURN_NOT_OK(
106-
ValidateFileConsistency(main_data_file_meta, blob_metas, blob_schema_->num_fields()));
107112

108-
results_.push_back(main_data_file_meta);
113+
if (main_data_file_meta != nullptr) {
114+
PAIMON_RETURN_NOT_OK(
115+
ValidateFileConsistency(main_data_file_meta, blob_metas, blob_schema_->num_fields()));
116+
results_.push_back(main_data_file_meta);
117+
}
109118
results_.insert(results_.end(), blob_metas.begin(), blob_metas.end());
110119

111120
current_writer_.reset();

src/paimon/core/io/rolling_blob_file_writer.h

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -38,7 +38,8 @@ namespace paimon {
3838
/// between them.
3939
///
4040
/// Multiple blob fields are supported. Each blob field is written to its own set of blob files
41-
/// independently via MultipleBlobFileWriter.
41+
/// independently via MultipleBlobFileWriter. For blob-only writes, the main writer factory may be
42+
/// nullptr and only blob files are produced.
4243
///
4344
/// <pre>
4445
/// For example,

src/paimon/core/schema/arrow_schema_validator.cpp

Lines changed: 34 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -55,10 +55,19 @@ Status ArrowSchemaValidator::ValidateSchema(const arrow::Schema& schema) {
5555

5656
Status ArrowSchemaValidator::ValidateSchemaWithFieldId(const arrow::Schema& schema) {
5757
PAIMON_RETURN_NOT_OK(ValidateSchema(schema));
58-
auto struct_type = arrow::struct_(schema.fields());
5958
std::set<int32_t> field_id_set;
60-
PAIMON_RETURN_NOT_OK(
61-
ValidateDataTypeWithFieldId(struct_type, /*key_value_metadata=*/nullptr, &field_id_set));
59+
for (const auto& field : schema.fields()) {
60+
PAIMON_ASSIGN_OR_RAISE(DataField data_field,
61+
DataField::ConvertArrowFieldToDataField(field));
62+
auto iter = field_id_set.find(data_field.Id());
63+
if (iter != field_id_set.end()) {
64+
return Status::Invalid(
65+
fmt::format("field id must be unique, duplicate field id {}", data_field.Id()));
66+
}
67+
field_id_set.insert(data_field.Id());
68+
PAIMON_RETURN_NOT_OK(ValidateDataTypeWithFieldId(field->type(), field->metadata(),
69+
/*allow_blob=*/true, &field_id_set));
70+
}
6271
return Status::OK();
6372
}
6473

@@ -93,7 +102,7 @@ Status ArrowSchemaValidator::ValidateNoWhitespaceOnlyFields(const arrow::FieldVe
93102

94103
Status ArrowSchemaValidator::ValidateDataTypeWithFieldId(
95104
const std::shared_ptr<arrow::DataType>& type,
96-
const std::shared_ptr<const arrow::KeyValueMetadata>& key_value_metadata,
105+
const std::shared_ptr<const arrow::KeyValueMetadata>& key_value_metadata, bool allow_blob,
97106
std::set<int32_t>* field_id_set) {
98107
const auto kind = type->id();
99108
switch (kind) {
@@ -114,7 +123,7 @@ Status ArrowSchemaValidator::ValidateDataTypeWithFieldId(
114123
const auto& value_field =
115124
arrow::internal::checked_cast<arrow::BaseListType*>(type.get())->value_field();
116125
PAIMON_RETURN_NOT_OK(ValidateDataTypeWithFieldId(
117-
value_field->type(), value_field->metadata(), field_id_set));
126+
value_field->type(), value_field->metadata(), /*allow_blob=*/false, field_id_set));
118127
break;
119128
}
120129
case arrow::Type::type::STRUCT: {
@@ -135,7 +144,7 @@ Status ArrowSchemaValidator::ValidateDataTypeWithFieldId(
135144
}
136145
field_id_set->insert(data_field.Id());
137146
PAIMON_RETURN_NOT_OK(ValidateDataTypeWithFieldId(
138-
sub_field->type(), sub_field->metadata(), field_id_set));
147+
sub_field->type(), sub_field->metadata(), /*allow_blob=*/false, field_id_set));
139148
}
140149
break;
141150
}
@@ -144,14 +153,17 @@ Status ArrowSchemaValidator::ValidateDataTypeWithFieldId(
144153
arrow::internal::checked_cast<arrow::MapType*>(type.get())->key_field();
145154
const auto& item_field =
146155
arrow::internal::checked_cast<arrow::MapType*>(type.get())->item_field();
147-
PAIMON_RETURN_NOT_OK(ValidateDataTypeWithFieldId(key_field->type(),
148-
key_field->metadata(), field_id_set));
149-
PAIMON_RETURN_NOT_OK(ValidateDataTypeWithFieldId(item_field->type(),
150-
item_field->metadata(), field_id_set));
156+
PAIMON_RETURN_NOT_OK(ValidateDataTypeWithFieldId(
157+
key_field->type(), key_field->metadata(), /*allow_blob=*/false, field_id_set));
158+
PAIMON_RETURN_NOT_OK(ValidateDataTypeWithFieldId(
159+
item_field->type(), item_field->metadata(), /*allow_blob=*/false, field_id_set));
151160
break;
152161
}
153162
case arrow::Type::type::LARGE_BINARY: {
154163
if (BlobUtils::IsBlobMetadata(key_value_metadata)) {
164+
if (!allow_blob) {
165+
return Status::Invalid("Blob field must be a top-level field.");
166+
}
155167
break;
156168
}
157169
[[fallthrough]];
@@ -164,6 +176,11 @@ Status ArrowSchemaValidator::ValidateDataTypeWithFieldId(
164176
}
165177

166178
Status ArrowSchemaValidator::ValidateField(const std::shared_ptr<arrow::Field>& field) {
179+
return ValidateField(field, /*allow_blob=*/true);
180+
}
181+
182+
Status ArrowSchemaValidator::ValidateField(const std::shared_ptr<arrow::Field>& field,
183+
bool allow_blob) {
167184
const auto kind = field->type()->id();
168185
switch (kind) {
169186
case arrow::Type::type::BOOL:
@@ -185,7 +202,7 @@ Status ArrowSchemaValidator::ValidateField(const std::shared_ptr<arrow::Field>&
185202
const auto& value_field =
186203
arrow::internal::checked_cast<const arrow::BaseListType&>(*field->type())
187204
.value_field();
188-
PAIMON_RETURN_NOT_OK(ValidateField(value_field));
205+
PAIMON_RETURN_NOT_OK(ValidateField(value_field, /*allow_blob=*/false));
189206
break;
190207
}
191208
case arrow::Type::type::STRUCT: {
@@ -202,7 +219,7 @@ Status ArrowSchemaValidator::ValidateField(const std::shared_ptr<arrow::Field>&
202219
arrow::FieldVector arrow_fields =
203220
arrow::internal::checked_cast<const arrow::StructType&>(*field->type()).fields();
204221
for (const auto& sub_field : arrow_fields) {
205-
PAIMON_RETURN_NOT_OK(ValidateField(sub_field));
222+
PAIMON_RETURN_NOT_OK(ValidateField(sub_field, /*allow_blob=*/false));
206223
}
207224
break;
208225
}
@@ -211,12 +228,15 @@ Status ArrowSchemaValidator::ValidateField(const std::shared_ptr<arrow::Field>&
211228
arrow::internal::checked_cast<const arrow::MapType&>(*field->type()).key_field();
212229
const auto& item_field =
213230
arrow::internal::checked_cast<const arrow::MapType&>(*field->type()).item_field();
214-
PAIMON_RETURN_NOT_OK(ValidateField(key_field));
215-
PAIMON_RETURN_NOT_OK(ValidateField(item_field));
231+
PAIMON_RETURN_NOT_OK(ValidateField(key_field, /*allow_blob=*/false));
232+
PAIMON_RETURN_NOT_OK(ValidateField(item_field, /*allow_blob=*/false));
216233
break;
217234
}
218235
case arrow::Type::type::LARGE_BINARY: {
219236
if (BlobUtils::IsBlobField(field)) {
237+
if (!allow_blob) {
238+
return Status::Invalid("Blob field must be a top-level field.");
239+
}
220240
break;
221241
}
222242
[[fallthrough]];

src/paimon/core/schema/arrow_schema_validator.h

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -55,7 +55,8 @@ class PAIMON_EXPORT ArrowSchemaValidator {
5555
private:
5656
static Status ValidateDataTypeWithFieldId(
5757
const std::shared_ptr<arrow::DataType>& type,
58-
const std::shared_ptr<const arrow::KeyValueMetadata>& key_value_metadata,
58+
const std::shared_ptr<const arrow::KeyValueMetadata>& key_value_metadata, bool allow_blob,
5959
std::set<int32_t>* field_id_set);
60+
static Status ValidateField(const std::shared_ptr<arrow::Field>& field, bool allow_blob);
6061
};
6162
} // namespace paimon

src/paimon/core/schema/arrow_schema_validator_test.cpp

Lines changed: 48 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -21,6 +21,7 @@
2121

2222
#include "arrow/type.h"
2323
#include "gtest/gtest.h"
24+
#include "paimon/common/data/blob_utils.h"
2425
#include "paimon/common/data/variant/variant_access_utils.h"
2526
#include "paimon/common/data/variant/variant_defs.h"
2627
#include "paimon/common/data/variant/variant_type_utils.h"
@@ -151,6 +152,51 @@ TEST(ArrowSchemaValidatorTest, TestInvalidDataType) {
151152
}
152153
}
153154

155+
TEST(ArrowSchemaValidatorTest, TestBlobFieldMustBeTopLevel) {
156+
{
157+
auto arrow_schema =
158+
arrow::schema(arrow::FieldVector({BlobUtils::ToArrowField("blob", true)}));
159+
ASSERT_OK(ArrowSchemaValidator::ValidateSchema(*arrow_schema));
160+
}
161+
{
162+
std::vector<DataField> fields = {DataField(0, BlobUtils::ToArrowField("blob", true))};
163+
auto arrow_schema = DataField::ConvertDataFieldsToArrowSchema(fields);
164+
ASSERT_OK(ArrowSchemaValidator::ValidateSchemaWithFieldId(*arrow_schema));
165+
}
166+
{
167+
auto nested_blob_field =
168+
arrow::field("nested", arrow::struct_({BlobUtils::ToArrowField("blob", true)}));
169+
auto arrow_schema = arrow::schema(arrow::FieldVector({nested_blob_field}));
170+
ASSERT_NOK_WITH_MSG(ArrowSchemaValidator::ValidateSchema(*arrow_schema),
171+
"Blob field must be a top-level field.");
172+
}
173+
{
174+
auto array_blob_field =
175+
arrow::field("array_blob", arrow::list(BlobUtils::ToArrowField("item", true)));
176+
auto arrow_schema = arrow::schema(arrow::FieldVector({array_blob_field}));
177+
ASSERT_NOK_WITH_MSG(ArrowSchemaValidator::ValidateSchema(*arrow_schema),
178+
"Blob field must be a top-level field.");
179+
}
180+
{
181+
auto map_blob_field = arrow::field(
182+
"map_blob",
183+
arrow::map(arrow::utf8(), arrow::struct_({BlobUtils::ToArrowField("blob", true)})));
184+
auto arrow_schema = arrow::schema(arrow::FieldVector({map_blob_field}));
185+
ASSERT_NOK_WITH_MSG(ArrowSchemaValidator::ValidateSchema(*arrow_schema),
186+
"Blob field must be a top-level field.");
187+
}
188+
{
189+
std::vector<DataField> nested_fields = {
190+
DataField(1, BlobUtils::ToArrowField("blob", true))};
191+
std::vector<DataField> fields = {DataField(
192+
0,
193+
arrow::field("nested", DataField::ConvertDataFieldsToArrowStructType(nested_fields)))};
194+
auto arrow_schema = DataField::ConvertDataFieldsToArrowSchema(fields);
195+
ASSERT_NOK_WITH_MSG(ArrowSchemaValidator::ValidateSchemaWithFieldId(*arrow_schema),
196+
"Blob field must be a top-level field.");
197+
}
198+
}
199+
154200
TEST(ArrowSchemaValidatorTest, ValidateDataTypeWithFieldId) {
155201
{
156202
std::vector<DataField> fields = {DataField(3, arrow::field("f3", arrow::float64())),
@@ -268,7 +314,8 @@ TEST(ArrowSchemaValidatorTest, ValidateDataTypeWithFieldId) {
268314
auto struct_type = DataField::ConvertDataFieldsToArrowStructType(fields);
269315
std::set<int32_t> field_id_set;
270316
ASSERT_NOK_WITH_MSG(ArrowSchemaValidator::ValidateDataTypeWithFieldId(
271-
struct_type, /*key_value_metadata=*/nullptr, &field_id_set),
317+
struct_type, /*key_value_metadata=*/nullptr,
318+
/*allow_blob=*/true, &field_id_set),
272319
"Unknown or unsupported arrow type: large_string");
273320
}
274321
}

0 commit comments

Comments
 (0)