Skip to content

Commit 31620d3

Browse files
committed
fix review
1 parent d471f04 commit 31620d3

4 files changed

Lines changed: 44 additions & 36 deletions

File tree

src/paimon/core/append/append_only_writer_test.cpp

Lines changed: 16 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -763,16 +763,29 @@ TEST_F(AppendOnlyWriterTest, TestWriteWithOnlyBlobField) {
763763

764764
ASSERT_OK(writer->Write(CreateStructBatch(schema, {blob_array})));
765765
ASSERT_OK_AND_ASSIGN(CommitIncrement inc, writer->PrepareCommit(/*wait_compaction=*/true));
766+
ASSERT_OK(writer->Close());
766767

767768
const auto& new_files = inc.GetNewFilesIncrement().NewFiles();
768769
ASSERT_EQ(new_files.size(), 1);
769770
ASSERT_TRUE(BlobUtils::IsBlobFile(new_files[0]->file_name));
770771
ASSERT_EQ(new_files[0]->row_count, 2);
771772
ASSERT_TRUE(new_files[0]->write_cols.has_value());
772773
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());
774+
std::string blob_file_path = path_factory->ToPath(new_files[0]->file_name);
775+
ASSERT_TRUE(options.GetFileSystem()->Exists(blob_file_path).value());
776+
777+
auto blob_reader = OpenFormatReader(blob_file_path, "blob");
778+
::ArrowSchema c_blob_schema;
779+
ASSERT_TRUE(arrow::ExportSchema(*schema, &c_blob_schema).ok());
780+
ASSERT_OK(blob_reader->SetReadSchema(&c_blob_schema, /*predicate=*/nullptr,
781+
/*selection_bitmap=*/std::nullopt));
782+
ASSERT_OK_AND_ASSIGN(auto actual_array, ReadResultCollector::CollectResult(blob_reader.get()));
783+
auto expected_struct_array =
784+
arrow::StructArray::Make({blob_array}, {blob_field->name()}).ValueOrDie();
785+
auto expected_array = std::make_shared<arrow::ChunkedArray>(expected_struct_array);
786+
ASSERT_TRUE(expected_array->Equals(actual_array)) << "Expected:\n"
787+
<< expected_array->ToString() << "\nActual:\n"
788+
<< actual_array->ToString();
776789
}
777790

778791
TEST_F(AppendOnlyWriterTest, TestWriteWithMultipleBlobFields) {

src/paimon/core/io/rolling_blob_file_writer.cpp

Lines changed: 18 additions & 26 deletions
Original file line numberDiff line numberDiff line change
@@ -16,6 +16,7 @@
1616

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

19+
#include <map>
1920
#include <memory>
2021
#include <string>
2122
#include <utility>
@@ -27,7 +28,6 @@
2728
#include "arrow/c/bridge.h"
2829
#include "arrow/c/helpers.h"
2930
#include "fmt/format.h"
30-
#include "fmt/ranges.h"
3131
#include "paimon/common/data/blob_utils.h"
3232
#include "paimon/common/utils/arrow/status_utils.h"
3333
#include "paimon/common/utils/scope_guard.h"
@@ -97,9 +97,6 @@ Status RollingBlobFileWriter::Write(::ArrowArray* record) {
9797
}
9898

9999
Status RollingBlobFileWriter::CloseCurrentWriter() {
100-
if (current_writer_ == nullptr && blob_writer_ == nullptr) {
101-
return Status::OK();
102-
}
103100
if (blob_writer_ == nullptr) {
104101
return Status::OK();
105102
}
@@ -111,8 +108,7 @@ Status RollingBlobFileWriter::CloseCurrentWriter() {
111108
CloseBlobWriter());
112109

113110
if (main_data_file_meta != nullptr) {
114-
PAIMON_RETURN_NOT_OK(
115-
ValidateFileConsistency(main_data_file_meta, blob_metas, blob_schema_->num_fields()));
111+
PAIMON_RETURN_NOT_OK(ValidateFileConsistency(main_data_file_meta, blob_metas));
116112
results_.push_back(main_data_file_meta);
117113
}
118114
results_.insert(results_.end(), blob_metas.begin(), blob_metas.end());
@@ -146,29 +142,25 @@ Result<std::vector<std::shared_ptr<DataFileMeta>>> RollingBlobFileWriter::CloseB
146142

147143
Status RollingBlobFileWriter::ValidateFileConsistency(
148144
const std::shared_ptr<DataFileMeta>& main_data_file_meta,
149-
const std::vector<std::shared_ptr<DataFileMeta>>& blob_tagged_metas, int32_t blob_field_count) {
150-
if (blob_tagged_metas.empty()) {
151-
return Status::OK();
152-
}
153-
// With multiple blob fields, each blob field produces its own set of files.
154-
// total_blob_row_count should be exactly main_row_count * blob_field_count.
155-
int64_t main_row_count = main_data_file_meta->row_count;
156-
int64_t expected_blob_row_count = main_row_count * blob_field_count;
157-
int64_t total_blob_row_count = 0;
145+
const std::vector<std::shared_ptr<DataFileMeta>>& blob_tagged_metas) {
146+
std::map<std::string, int64_t> blob_field_row_counts;
158147
for (const auto& blob_tagged_meta : blob_tagged_metas) {
159-
total_blob_row_count += blob_tagged_meta->row_count;
148+
if (!blob_tagged_meta->write_cols || blob_tagged_meta->write_cols->empty()) {
149+
return Status::Invalid(
150+
fmt::format("This is a bug: Blob file {} must contain a write column.",
151+
blob_tagged_meta->file_name));
152+
}
153+
blob_field_row_counts[blob_tagged_meta->write_cols->at(0)] += blob_tagged_meta->row_count;
160154
}
161-
if (total_blob_row_count != expected_blob_row_count) {
162-
std::vector<std::string> blob_file_names;
163-
for (const auto& blob_tagged_meta : blob_tagged_metas) {
164-
blob_file_names.push_back(blob_tagged_meta->file_name);
155+
156+
int64_t main_row_count = main_data_file_meta->row_count;
157+
for (const auto& [field_name, row_count] : blob_field_row_counts) {
158+
if (row_count != main_row_count) {
159+
return Status::Invalid(fmt::format(
160+
"This is a bug: The row count of main file and blob file does not match. Main "
161+
"file: {} (row count: {}), blob field name: {} (row count: {})",
162+
main_data_file_meta->file_name, main_row_count, field_name, row_count));
165163
}
166-
return Status::Invalid(fmt::format(
167-
"This is a bug: The row count of main file and blob files does not match. "
168-
"Main file: {} (row count: {}), blob field count: {}, "
169-
"expected blob row count: {}, blob files: {} (actual total row count: {})",
170-
main_data_file_meta->file_name, main_row_count, blob_field_count,
171-
expected_blob_row_count, fmt::join(blob_file_names, ", "), total_blob_row_count));
172164
}
173165
return Status::OK();
174166
}

src/paimon/core/io/rolling_blob_file_writer.h

Lines changed: 1 addition & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -77,8 +77,7 @@ class RollingBlobFileWriter
7777
private:
7878
static Status ValidateFileConsistency(
7979
const std::shared_ptr<DataFileMeta>& main_data_file_meta,
80-
const std::vector<std::shared_ptr<DataFileMeta>>& blob_tagged_metas,
81-
int32_t blob_field_count);
80+
const std::vector<std::shared_ptr<DataFileMeta>>& blob_tagged_metas);
8281

8382
Status CloseCurrentWriter();
8483

src/paimon/core/io/rolling_blob_file_writer_test.cpp

Lines changed: 9 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -82,11 +82,15 @@ TEST_F(RollingBlobFileWriterTest, ValidateFileConsistency) {
8282
/*delete_row_count=*/0, /*embedded_index=*/nullptr, FileSource::Append(),
8383
/*value_stats_cols=*/std::nullopt, /*external_path=*/std::nullopt, /*first_row_id=*/3,
8484
/*write_cols=*/std::vector<std::string>({"blob"}));
85-
ASSERT_OK(RollingBlobFileWriter::ValidateFileConsistency(file_meta1, {file_meta2, file_meta3},
86-
/*blob_field_count=*/1));
87-
ASSERT_NOK_WITH_MSG(RollingBlobFileWriter::ValidateFileConsistency(file_meta1, {file_meta2},
88-
/*blob_field_count=*/2),
89-
"This is a bug: The row count of main file and blob files does not match.");
85+
ASSERT_OK(RollingBlobFileWriter::ValidateFileConsistency(file_meta1, {file_meta2, file_meta3}));
86+
ASSERT_NOK_WITH_MSG(RollingBlobFileWriter::ValidateFileConsistency(file_meta1, {file_meta2}),
87+
"This is a bug: The row count of main file and blob file does not match.");
88+
89+
file_meta2->write_cols = std::vector<std::string>({"blob1"});
90+
file_meta3->write_cols = std::vector<std::string>({"blob2"});
91+
ASSERT_NOK_WITH_MSG(
92+
RollingBlobFileWriter::ValidateFileConsistency(file_meta1, {file_meta2, file_meta3}),
93+
"This is a bug: The row count of main file and blob file does not match.");
9094
}
9195

9296
} // namespace paimon::test

0 commit comments

Comments
 (0)