1616
1717#include " paimon/core/io/rolling_blob_file_writer.h"
1818
19+ #include < map>
1920#include < memory>
2021#include < string>
2122#include < utility>
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
9999Status 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
147143Status 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}
0 commit comments