Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
9 changes: 4 additions & 5 deletions src/paimon/common/data/blob_utils.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -82,13 +82,12 @@ Result<BlobUtils::SeparatedStructArrays> BlobUtils::SeparateBlobArray(
return Status::Invalid(
"SeparateBlobArray expects at least one non-inline blob field, but got none.");
}
if (main_fields.empty()) {
return Status::Invalid("SeparateBlobArray expects at least one main field, but got none.");
}

SeparatedStructArrays result;
PAIMON_ASSIGN_OR_RAISE_FROM_ARROW(result.main_array,
arrow::StructArray::Make(main_arrays, main_fields));
if (!main_fields.empty()) {
PAIMON_ASSIGN_OR_RAISE_FROM_ARROW(result.main_array,
arrow::StructArray::Make(main_arrays, main_fields));
}
PAIMON_ASSIGN_OR_RAISE_FROM_ARROW(result.blob_array,
arrow::StructArray::Make(blob_arrays, blob_fields));
return result;
Expand Down
3 changes: 2 additions & 1 deletion src/paimon/common/data/blob_utils.h
Original file line number Diff line number Diff line change
Expand Up @@ -51,7 +51,8 @@ class PAIMON_EXPORT BlobUtils {
};

struct SeparatedStructArrays {
/// Non-blob fields (includes inline blob fields when inline_fields is provided)
/// Non-blob fields (includes inline blob fields when inline_fields is provided).
/// nullptr when all fields are stored in blob files.
std::shared_ptr<arrow::StructArray> main_array;
/// Blob fields that go to separate .blob files
std::shared_ptr<arrow::StructArray> blob_array;
Expand Down
8 changes: 5 additions & 3 deletions src/paimon/common/data/blob_utils_test.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -182,11 +182,13 @@ TEST_F(BlobUtilsTest, SeparateBlobArray) {
BlobUtils::SeparateBlobArray(struct_array, /*inline_fields=*/{"f2_blob"}),
"SeparateBlobArray expects at least one non-inline blob field, but got none.");

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

TEST_F(BlobUtilsTest, SeparateBlobArrayWithPartialInline) {
Expand Down
10 changes: 7 additions & 3 deletions src/paimon/core/append/append_only_writer.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -236,10 +236,14 @@ AppendOnlyWriter::RollingFileWriterResult AppendOnlyWriter::CreateRollingBlobWri
options_.GetBlobTargetFileSize(), single_blob_file_writer_factory);
};

WriterFactory main_writer_factory;
if (schemas.main_schema->num_fields() > 0) {
main_writer_factory =
GetDataFileWriterFactory(schemas.main_schema, schemas.main_schema->field_names());
}
return std::make_unique<RollingBlobFileWriter>(
options_.GetTargetFileSize(/*has_primary_key=*/false),
GetDataFileWriterFactory(schemas.main_schema, schemas.main_schema->field_names()),
blob_schema, blob_writer_creator, arrow::struct_(write_schema_->fields()), inline_fields);
options_.GetTargetFileSize(/*has_primary_key=*/false), main_writer_factory, blob_schema,
blob_writer_creator, arrow::struct_(write_schema_->fields()), inline_fields);
}

Status AppendOnlyWriter::Sync() {
Expand Down
34 changes: 34 additions & 0 deletions src/paimon/core/append/append_only_writer_test.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -741,6 +741,40 @@ TEST_F(AppendOnlyWriterTest, TestWriteWithSingleBlobField) {
ASSERT_OK(writer->Close());
}

TEST_F(AppendOnlyWriterTest, TestWriteWithOnlyBlobField) {
auto options =
CreateOptions({{Options::FILE_FORMAT, "orc"}, {Options::MANIFEST_FORMAT, "orc"}});
auto dir = UniqueTestDirectory::Create();
ASSERT_TRUE(dir);
auto path_factory = CreatePathFactory(dir->Str(), "orc", options);

auto blob_field = BlobUtils::ToArrowField("blob", false);
auto schema = arrow::schema({blob_field});
ASSERT_OK_AND_ASSIGN(auto writer,
CreateAppendOnlyWriter(options, /*schema_id=*/0, schema,
/*write_cols=*/std::vector<std::string>{"blob"},
/*max_sequence_number=*/-1, path_factory,
compact_manager_, memory_pool_));

arrow::LargeBinaryBuilder blob_builder;
ASSERT_TRUE(blob_builder.Append("a", 1).ok());
ASSERT_TRUE(blob_builder.Append("bb", 2).ok());
auto blob_array = blob_builder.Finish().ValueOrDie();

ASSERT_OK(writer->Write(CreateStructBatch(schema, {blob_array})));
ASSERT_OK_AND_ASSIGN(CommitIncrement inc, writer->PrepareCommit(/*wait_compaction=*/true));

const auto& new_files = inc.GetNewFilesIncrement().NewFiles();
ASSERT_EQ(new_files.size(), 1);
ASSERT_TRUE(BlobUtils::IsBlobFile(new_files[0]->file_name));
ASSERT_EQ(new_files[0]->row_count, 2);
ASSERT_TRUE(new_files[0]->write_cols.has_value());
ASSERT_EQ(new_files[0]->write_cols.value(), std::vector<std::string>({"blob"}));
ASSERT_TRUE(
options.GetFileSystem()->Exists(path_factory->ToPath(new_files[0]->file_name)).value());
ASSERT_OK(writer->Close());
}

TEST_F(AppendOnlyWriterTest, TestWriteWithMultipleBlobFields) {
auto options =
CreateOptions({{Options::FILE_FORMAT, "orc"}, {Options::MANIFEST_FORMAT, "orc"}});
Expand Down
39 changes: 24 additions & 15 deletions src/paimon/core/io/rolling_blob_file_writer.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -55,7 +55,7 @@ RollingBlobFileWriter::RollingBlobFileWriter(
Status RollingBlobFileWriter::Write(::ArrowArray* record) {
ScopeGuard guard([this]() -> void { this->Abort(); });
// Open the current writer if write the first record or roll over happen before.
if (PAIMON_UNLIKELY(current_writer_ == nullptr)) {
if (writer_factory_ != nullptr && PAIMON_UNLIKELY(current_writer_ == nullptr)) {
PAIMON_RETURN_NOT_OK(OpenCurrentWriter());
}
if (PAIMON_UNLIKELY(blob_writer_ == nullptr)) {
Expand All @@ -69,12 +69,14 @@ Status RollingBlobFileWriter::Write(::ArrowArray* record) {
PAIMON_ASSIGN_OR_RAISE(BlobUtils::SeparatedStructArrays separated_arrays,
BlobUtils::SeparateBlobArray(struct_array, inline_fields_));
// Write main (non-blob) data
::ArrowArray c_main_array;
PAIMON_RETURN_NOT_OK_FROM_ARROW(
arrow::ExportArray(*separated_arrays.main_array, &c_main_array));
ScopeGuard array_lifecycle_guard(
[&c_main_array]() -> void { ArrowArrayRelease(&c_main_array); });
PAIMON_RETURN_NOT_OK(current_writer_->Write(&c_main_array));
if (current_writer_ != nullptr) {
::ArrowArray c_main_array;
PAIMON_RETURN_NOT_OK_FROM_ARROW(
arrow::ExportArray(*separated_arrays.main_array, &c_main_array));
ScopeGuard array_lifecycle_guard(
[&c_main_array]() -> void { ArrowArrayRelease(&c_main_array); });
PAIMON_RETURN_NOT_OK(current_writer_->Write(&c_main_array));
}

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

record_count_ += record_count;
PAIMON_ASSIGN_OR_RAISE(bool need_rolling_file, NeedRollingFile());
if (need_rolling_file) {
PAIMON_RETURN_NOT_OK(CloseCurrentWriter());
if (current_writer_ != nullptr) {
PAIMON_ASSIGN_OR_RAISE(bool need_rolling_file, NeedRollingFile());
if (need_rolling_file) {
PAIMON_RETURN_NOT_OK(CloseCurrentWriter());
}
}
guard.Release();
return Status::OK();
}

Status RollingBlobFileWriter::CloseCurrentWriter() {
if (current_writer_ == nullptr) {
if (current_writer_ == nullptr && blob_writer_ == nullptr) {
return Status::OK();
}
if (blob_writer_ == nullptr) {
return Status::OK();
}
PAIMON_ASSIGN_OR_RAISE(std::shared_ptr<DataFileMeta> main_data_file_meta, CloseMainWriter());
std::shared_ptr<DataFileMeta> main_data_file_meta;
if (current_writer_ != nullptr) {
PAIMON_ASSIGN_OR_RAISE(main_data_file_meta, CloseMainWriter());
}
PAIMON_ASSIGN_OR_RAISE(std::vector<std::shared_ptr<DataFileMeta>> blob_metas,
CloseBlobWriter());
PAIMON_RETURN_NOT_OK(
ValidateFileConsistency(main_data_file_meta, blob_metas, blob_schema_->num_fields()));

results_.push_back(main_data_file_meta);
if (main_data_file_meta != nullptr) {
PAIMON_RETURN_NOT_OK(
ValidateFileConsistency(main_data_file_meta, blob_metas, blob_schema_->num_fields()));
results_.push_back(main_data_file_meta);
}
results_.insert(results_.end(), blob_metas.begin(), blob_metas.end());

current_writer_.reset();
Expand Down
3 changes: 2 additions & 1 deletion src/paimon/core/io/rolling_blob_file_writer.h
Original file line number Diff line number Diff line change
Expand Up @@ -38,7 +38,8 @@ namespace paimon {
/// between them.
///
/// Multiple blob fields are supported. Each blob field is written to its own set of blob files
/// independently via MultipleBlobFileWriter.
/// independently via MultipleBlobFileWriter. For blob-only writes, the main writer factory may be
/// nullptr and only blob files are produced.
///
/// <pre>
/// For example,
Expand Down
48 changes: 34 additions & 14 deletions src/paimon/core/schema/arrow_schema_validator.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -53,10 +53,19 @@ Status ArrowSchemaValidator::ValidateSchema(const arrow::Schema& schema) {

Status ArrowSchemaValidator::ValidateSchemaWithFieldId(const arrow::Schema& schema) {
PAIMON_RETURN_NOT_OK(ValidateSchema(schema));
auto struct_type = arrow::struct_(schema.fields());
std::set<int32_t> field_id_set;
PAIMON_RETURN_NOT_OK(
ValidateDataTypeWithFieldId(struct_type, /*key_value_metadata=*/nullptr, &field_id_set));
for (const auto& field : schema.fields()) {
PAIMON_ASSIGN_OR_RAISE(DataField data_field,
DataField::ConvertArrowFieldToDataField(field));
auto iter = field_id_set.find(data_field.Id());
if (iter != field_id_set.end()) {
return Status::Invalid(
fmt::format("field id must be unique, duplicate field id {}", data_field.Id()));
}
field_id_set.insert(data_field.Id());
PAIMON_RETURN_NOT_OK(ValidateDataTypeWithFieldId(field->type(), field->metadata(),
/*allow_blob=*/true, &field_id_set));
}
return Status::OK();
}

Expand Down Expand Up @@ -91,7 +100,7 @@ Status ArrowSchemaValidator::ValidateNoWhitespaceOnlyFields(const arrow::FieldVe

Status ArrowSchemaValidator::ValidateDataTypeWithFieldId(
const std::shared_ptr<arrow::DataType>& type,
const std::shared_ptr<const arrow::KeyValueMetadata>& key_value_metadata,
const std::shared_ptr<const arrow::KeyValueMetadata>& key_value_metadata, bool allow_blob,
std::set<int32_t>* field_id_set) {
const auto kind = type->id();
switch (kind) {
Expand All @@ -112,7 +121,7 @@ Status ArrowSchemaValidator::ValidateDataTypeWithFieldId(
const auto& value_field =
arrow::internal::checked_cast<arrow::BaseListType*>(type.get())->value_field();
PAIMON_RETURN_NOT_OK(ValidateDataTypeWithFieldId(
value_field->type(), value_field->metadata(), field_id_set));
value_field->type(), value_field->metadata(), /*allow_blob=*/false, field_id_set));
break;
}
case arrow::Type::type::STRUCT: {
Expand All @@ -128,7 +137,7 @@ Status ArrowSchemaValidator::ValidateDataTypeWithFieldId(
}
field_id_set->insert(data_field.Id());
PAIMON_RETURN_NOT_OK(ValidateDataTypeWithFieldId(
sub_field->type(), sub_field->metadata(), field_id_set));
sub_field->type(), sub_field->metadata(), /*allow_blob=*/false, field_id_set));
}
break;
}
Expand All @@ -137,14 +146,17 @@ Status ArrowSchemaValidator::ValidateDataTypeWithFieldId(
arrow::internal::checked_cast<arrow::MapType*>(type.get())->key_field();
const auto& item_field =
arrow::internal::checked_cast<arrow::MapType*>(type.get())->item_field();
PAIMON_RETURN_NOT_OK(ValidateDataTypeWithFieldId(key_field->type(),
key_field->metadata(), field_id_set));
PAIMON_RETURN_NOT_OK(ValidateDataTypeWithFieldId(item_field->type(),
item_field->metadata(), field_id_set));
PAIMON_RETURN_NOT_OK(ValidateDataTypeWithFieldId(
key_field->type(), key_field->metadata(), /*allow_blob=*/false, field_id_set));
PAIMON_RETURN_NOT_OK(ValidateDataTypeWithFieldId(
item_field->type(), item_field->metadata(), /*allow_blob=*/false, field_id_set));
break;
}
case arrow::Type::type::LARGE_BINARY: {
if (BlobUtils::IsBlobMetadata(key_value_metadata)) {
if (!allow_blob) {
return Status::Invalid("Blob field must be a top-level field.");
}
break;
}
[[fallthrough]];
Expand All @@ -157,6 +169,11 @@ Status ArrowSchemaValidator::ValidateDataTypeWithFieldId(
}

Status ArrowSchemaValidator::ValidateField(const std::shared_ptr<arrow::Field>& field) {
return ValidateField(field, /*allow_blob=*/true);
}

Status ArrowSchemaValidator::ValidateField(const std::shared_ptr<arrow::Field>& field,
bool allow_blob) {
const auto kind = field->type()->id();
switch (kind) {
case arrow::Type::type::BOOL:
Expand All @@ -178,14 +195,14 @@ Status ArrowSchemaValidator::ValidateField(const std::shared_ptr<arrow::Field>&
const auto& value_field =
arrow::internal::checked_cast<const arrow::BaseListType&>(*field->type())
.value_field();
PAIMON_RETURN_NOT_OK(ValidateField(value_field));
PAIMON_RETURN_NOT_OK(ValidateField(value_field, /*allow_blob=*/false));
break;
}
case arrow::Type::type::STRUCT: {
arrow::FieldVector arrow_fields =
arrow::internal::checked_cast<const arrow::StructType&>(*field->type()).fields();
for (const auto& sub_field : arrow_fields) {
PAIMON_RETURN_NOT_OK(ValidateField(sub_field));
PAIMON_RETURN_NOT_OK(ValidateField(sub_field, /*allow_blob=*/false));
}
break;
}
Expand All @@ -194,12 +211,15 @@ Status ArrowSchemaValidator::ValidateField(const std::shared_ptr<arrow::Field>&
arrow::internal::checked_cast<const arrow::MapType&>(*field->type()).key_field();
const auto& item_field =
arrow::internal::checked_cast<const arrow::MapType&>(*field->type()).item_field();
PAIMON_RETURN_NOT_OK(ValidateField(key_field));
PAIMON_RETURN_NOT_OK(ValidateField(item_field));
PAIMON_RETURN_NOT_OK(ValidateField(key_field, /*allow_blob=*/false));
PAIMON_RETURN_NOT_OK(ValidateField(item_field, /*allow_blob=*/false));
break;
}
case arrow::Type::type::LARGE_BINARY: {
if (BlobUtils::IsBlobField(field)) {
if (!allow_blob) {
return Status::Invalid("Blob field must be a top-level field.");
}
break;
}
[[fallthrough]];
Expand Down
3 changes: 2 additions & 1 deletion src/paimon/core/schema/arrow_schema_validator.h
Original file line number Diff line number Diff line change
Expand Up @@ -55,7 +55,8 @@ class PAIMON_EXPORT ArrowSchemaValidator {
private:
static Status ValidateDataTypeWithFieldId(
const std::shared_ptr<arrow::DataType>& type,
const std::shared_ptr<const arrow::KeyValueMetadata>& key_value_metadata,
const std::shared_ptr<const arrow::KeyValueMetadata>& key_value_metadata, bool allow_blob,
std::set<int32_t>* field_id_set);
static Status ValidateField(const std::shared_ptr<arrow::Field>& field, bool allow_blob);
};
} // namespace paimon
49 changes: 48 additions & 1 deletion src/paimon/core/schema/arrow_schema_validator_test.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,7 @@

#include "arrow/type.h"
#include "gtest/gtest.h"
#include "paimon/common/data/blob_utils.h"
#include "paimon/common/types/data_field.h"
#include "paimon/common/utils/date_time_utils.h"
#include "paimon/testing/utils/testharness.h"
Expand Down Expand Up @@ -148,6 +149,51 @@ TEST(ArrowSchemaValidatorTest, TestInvalidDataType) {
}
}

TEST(ArrowSchemaValidatorTest, TestBlobFieldMustBeTopLevel) {
{
auto arrow_schema =
arrow::schema(arrow::FieldVector({BlobUtils::ToArrowField("blob", true)}));
ASSERT_OK(ArrowSchemaValidator::ValidateSchema(*arrow_schema));
}
{
std::vector<DataField> fields = {DataField(0, BlobUtils::ToArrowField("blob", true))};
auto arrow_schema = DataField::ConvertDataFieldsToArrowSchema(fields);
ASSERT_OK(ArrowSchemaValidator::ValidateSchemaWithFieldId(*arrow_schema));
}
{
auto nested_blob_field =
arrow::field("nested", arrow::struct_({BlobUtils::ToArrowField("blob", true)}));
auto arrow_schema = arrow::schema(arrow::FieldVector({nested_blob_field}));
ASSERT_NOK_WITH_MSG(ArrowSchemaValidator::ValidateSchema(*arrow_schema),
"Blob field must be a top-level field.");
}
{
auto array_blob_field =
arrow::field("array_blob", arrow::list(BlobUtils::ToArrowField("item", true)));
auto arrow_schema = arrow::schema(arrow::FieldVector({array_blob_field}));
ASSERT_NOK_WITH_MSG(ArrowSchemaValidator::ValidateSchema(*arrow_schema),
"Blob field must be a top-level field.");
}
{
auto map_blob_field = arrow::field(
"map_blob",
arrow::map(arrow::utf8(), arrow::struct_({BlobUtils::ToArrowField("blob", true)})));
auto arrow_schema = arrow::schema(arrow::FieldVector({map_blob_field}));
ASSERT_NOK_WITH_MSG(ArrowSchemaValidator::ValidateSchema(*arrow_schema),
"Blob field must be a top-level field.");
}
{
std::vector<DataField> nested_fields = {
DataField(1, BlobUtils::ToArrowField("blob", true))};
std::vector<DataField> fields = {DataField(
0,
arrow::field("nested", DataField::ConvertDataFieldsToArrowStructType(nested_fields)))};
auto arrow_schema = DataField::ConvertDataFieldsToArrowSchema(fields);
ASSERT_NOK_WITH_MSG(ArrowSchemaValidator::ValidateSchemaWithFieldId(*arrow_schema),
"Blob field must be a top-level field.");
}
}

TEST(ArrowSchemaValidatorTest, ValidateDataTypeWithFieldId) {
{
std::vector<DataField> fields = {DataField(3, arrow::field("f3", arrow::float64())),
Expand Down Expand Up @@ -265,7 +311,8 @@ TEST(ArrowSchemaValidatorTest, ValidateDataTypeWithFieldId) {
auto struct_type = DataField::ConvertDataFieldsToArrowStructType(fields);
std::set<int32_t> field_id_set;
ASSERT_NOK_WITH_MSG(ArrowSchemaValidator::ValidateDataTypeWithFieldId(
struct_type, /*key_value_metadata=*/nullptr, &field_id_set),
struct_type, /*key_value_metadata=*/nullptr,
/*allow_blob=*/true, &field_id_set),
"Unknown or unsupported arrow type: large_string");
}
}
Expand Down
Loading
Loading