Skip to content
Merged
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 @@ -238,10 +238,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
47 changes: 47 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,53 @@ 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));
ASSERT_OK(writer->Close());

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"}));
std::string blob_file_path = path_factory->ToPath(new_files[0]->file_name);
ASSERT_TRUE(options.GetFileSystem()->Exists(blob_file_path).value());

auto blob_reader = OpenFormatReader(blob_file_path, "blob");
::ArrowSchema c_blob_schema;
ASSERT_TRUE(arrow::ExportSchema(*schema, &c_blob_schema).ok());
ASSERT_OK(blob_reader->SetReadSchema(&c_blob_schema, /*predicate=*/nullptr,
/*selection_bitmap=*/std::nullopt));
ASSERT_OK_AND_ASSIGN(auto actual_array, ReadResultCollector::CollectResult(blob_reader.get()));
auto expected_struct_array =
arrow::StructArray::Make({blob_array}, {blob_field->name()}).ValueOrDie();
auto expected_array = std::make_shared<arrow::ChunkedArray>(expected_struct_array);
ASSERT_TRUE(expected_array->Equals(actual_array)) << "Expected:\n"
<< expected_array->ToString() << "\nActual:\n"
<< actual_array->ToString();
}

TEST_F(AppendOnlyWriterTest, TestWriteWithMultipleBlobFields) {
auto options =
CreateOptions({{Options::FILE_FORMAT, "orc"}, {Options::MANIFEST_FORMAT, "orc"}});
Expand Down
77 changes: 39 additions & 38 deletions src/paimon/core/io/rolling_blob_file_writer.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@

#include "paimon/core/io/rolling_blob_file_writer.h"

#include <map>
#include <memory>
#include <string>
#include <utility>
Expand All @@ -27,7 +28,6 @@
#include "arrow/c/bridge.h"
#include "arrow/c/helpers.h"
#include "fmt/format.h"
#include "fmt/ranges.h"
#include "paimon/common/data/blob_utils.h"
#include "paimon/common/utils/arrow/status_utils.h"
#include "paimon/common/utils/scope_guard.h"
Expand Down 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,31 @@ 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) {
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));
results_.push_back(main_data_file_meta);
}
results_.insert(results_.end(), blob_metas.begin(), blob_metas.end());

current_writer_.reset();
Expand Down Expand Up @@ -137,29 +142,25 @@ Result<std::vector<std::shared_ptr<DataFileMeta>>> RollingBlobFileWriter::CloseB

Status RollingBlobFileWriter::ValidateFileConsistency(
const std::shared_ptr<DataFileMeta>& main_data_file_meta,
const std::vector<std::shared_ptr<DataFileMeta>>& blob_tagged_metas, int32_t blob_field_count) {
if (blob_tagged_metas.empty()) {
return Status::OK();
}
// With multiple blob fields, each blob field produces its own set of files.
// total_blob_row_count should be exactly main_row_count * blob_field_count.
int64_t main_row_count = main_data_file_meta->row_count;
int64_t expected_blob_row_count = main_row_count * blob_field_count;
int64_t total_blob_row_count = 0;
const std::vector<std::shared_ptr<DataFileMeta>>& blob_tagged_metas) {
std::map<std::string, int64_t> blob_field_row_counts;
for (const auto& blob_tagged_meta : blob_tagged_metas) {
total_blob_row_count += blob_tagged_meta->row_count;
if (!blob_tagged_meta->write_cols || blob_tagged_meta->write_cols->empty()) {
return Status::Invalid(
fmt::format("This is a bug: Blob file {} must contain a write column.",
blob_tagged_meta->file_name));
}
blob_field_row_counts[blob_tagged_meta->write_cols->at(0)] += blob_tagged_meta->row_count;
}
if (total_blob_row_count != expected_blob_row_count) {
std::vector<std::string> blob_file_names;
for (const auto& blob_tagged_meta : blob_tagged_metas) {
blob_file_names.push_back(blob_tagged_meta->file_name);

int64_t main_row_count = main_data_file_meta->row_count;
for (const auto& [field_name, row_count] : blob_field_row_counts) {
if (row_count != main_row_count) {
return Status::Invalid(fmt::format(
"This is a bug: The row count of main file and blob file does not match. Main "
"file: {} (row count: {}), blob field name: {} (row count: {})",
main_data_file_meta->file_name, main_row_count, field_name, row_count));
}
return Status::Invalid(fmt::format(
"This is a bug: The row count of main file and blob files does not match. "
"Main file: {} (row count: {}), blob field count: {}, "
"expected blob row count: {}, blob files: {} (actual total row count: {})",
main_data_file_meta->file_name, main_row_count, blob_field_count,
expected_blob_row_count, fmt::join(blob_file_names, ", "), total_blob_row_count));
}
return Status::OK();
}
Expand Down
6 changes: 3 additions & 3 deletions 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 Expand Up @@ -76,8 +77,7 @@ class RollingBlobFileWriter
private:
static Status ValidateFileConsistency(
const std::shared_ptr<DataFileMeta>& main_data_file_meta,
const std::vector<std::shared_ptr<DataFileMeta>>& blob_tagged_metas,
int32_t blob_field_count);
const std::vector<std::shared_ptr<DataFileMeta>>& blob_tagged_metas);

Status CloseCurrentWriter();

Expand Down
14 changes: 9 additions & 5 deletions src/paimon/core/io/rolling_blob_file_writer_test.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -82,11 +82,15 @@ TEST_F(RollingBlobFileWriterTest, ValidateFileConsistency) {
/*delete_row_count=*/0, /*embedded_index=*/nullptr, FileSource::Append(),
/*value_stats_cols=*/std::nullopt, /*external_path=*/std::nullopt, /*first_row_id=*/3,
/*write_cols=*/std::vector<std::string>({"blob"}));
ASSERT_OK(RollingBlobFileWriter::ValidateFileConsistency(file_meta1, {file_meta2, file_meta3},
/*blob_field_count=*/1));
ASSERT_NOK_WITH_MSG(RollingBlobFileWriter::ValidateFileConsistency(file_meta1, {file_meta2},
/*blob_field_count=*/2),
"This is a bug: The row count of main file and blob files does not match.");
ASSERT_OK(RollingBlobFileWriter::ValidateFileConsistency(file_meta1, {file_meta2, file_meta3}));
ASSERT_NOK_WITH_MSG(RollingBlobFileWriter::ValidateFileConsistency(file_meta1, {file_meta2}),
"This is a bug: The row count of main file and blob file does not match.");

file_meta2->write_cols = std::vector<std::string>({"blob1"});
file_meta3->write_cols = std::vector<std::string>({"blob2"});
ASSERT_NOK_WITH_MSG(
RollingBlobFileWriter::ValidateFileConsistency(file_meta1, {file_meta2, file_meta3}),
"This is a bug: The row count of main file and blob file does not match.");
}

} // namespace paimon::test
Loading
Loading