lucasfang commented on code in PR #395:
URL: https://github.com/apache/paimon-cpp/pull/395#discussion_r4128714562


##########
src/paimon/common/reader/data_evolution_file_reader.cpp:
##########
@@ -158,36 +171,39 @@ Result<std::shared_ptr<arrow::Array>> 
DataEvolutionFileReader::NextBatchForSingl
         PAIMON_ASSIGN_OR_RAISE(arrow::ArrayVector selected_array_vec,
                                
ReaderUtils::GenerateFilteredArrayVector(src_array, bitmap));
         for (const auto& selected_array : selected_array_vec) {
-            if (total_array_length + selected_array->length() > 
read_batch_size_) {
-                // need truncate current array to align read_batch_size_
-                int64_t truncated_length = read_batch_size_ - 
total_array_length;
-                if (truncated_length == 0) {
-                    // total_array_length equals to read_batch_size_, all 
selected_array left will
-                    // be added to cached_array_vec_
-                    cached_array_vec_[reader_idx].push_back(selected_array);
-                } else {
-                    concat_array_vec.push_back(selected_array->Slice(0, 
truncated_length));
-                    cached_array_vec_[reader_idx].push_back(
-                        selected_array->Slice(truncated_length));
-                    total_array_length += truncated_length;
-                }
-            } else {
-                concat_array_vec.push_back(selected_array);
-                total_array_length += selected_array->length();
+            if (selected_array->length() > 0) {
+                cached_array_vec_[reader_idx].push_back(selected_array);
             }
         }
     }
-    if (concat_array_vec.empty()) {
-        return std::shared_ptr<arrow::Array>();
+    return true;
+}
+
+Result<std::shared_ptr<arrow::Array>> DataEvolutionFileReader::TakeCachedArray(
+    size_t reader_idx, int64_t array_length) {
+    assert(array_length > 0);
+    assert(CalculateCachedArrayLength(reader_idx) >= array_length);
+    arrow::ArrayVector selected_array_vec;
+    int64_t remaining_length = array_length;
+    auto& cached_array_vec = cached_array_vec_[reader_idx];
+    while (remaining_length > 0) {
+        const auto& cached_array = cached_array_vec.front();
+        if (cached_array->length() <= remaining_length) {
+            selected_array_vec.push_back(cached_array);
+            remaining_length -= cached_array->length();
+            cached_array_vec.erase(cached_array_vec.begin());
+        } else {
+            selected_array_vec.push_back(cached_array->Slice(0, 
remaining_length));
+            cached_array_vec.front() = cached_array->Slice(remaining_length);
+            remaining_length = 0;
+        }
     }
-    if (concat_array_vec.size() == 1 && concat_array_vec[0]->offset() == 0) {
-        // Avoid data copy when the array is already normalized.
-        return concat_array_vec[0];
+    if (selected_array_vec.size() == 1) {
+        return ArrowUtils::NormalizeArrayOffsets(selected_array_vec[0], 
arrow_pool_.get());
     }
     PAIMON_ASSIGN_OR_RAISE_FROM_ARROW(std::shared_ptr<arrow::Array> 
concat_array,
-                                      arrow::Concatenate(concat_array_vec, 
arrow_pool_.get()));
-    assert(concat_array->length() == total_array_length);
-    assert(concat_array->length() <= read_batch_size_);
+                                      arrow::Concatenate(selected_array_vec, 
arrow_pool_.get()));
+    assert(concat_array->length() == array_length);

Review Comment:
   same



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to