lxy-9602 commented on code in PR #269:
URL: https://github.com/apache/paimon-cpp/pull/269#discussion_r3915112778


##########
src/paimon/core/io/key_value_data_file_record_reader.cpp:
##########
@@ -107,26 +110,34 @@ Result<std::unique_ptr<KeyValueRecordReader::Iterator>> 
KeyValueDataFileRecordRe
         return Status::Invalid("cannot cast data batch to StructArray");
     }
     auto data_batch = checked_pointer_cast<arrow::StructArray>(arrow_array);
-    if (data_batch->num_fields() < 
SpecialFields::KEY_VALUE_SPECIAL_FIELD_COUNT) {
-        return Status::Invalid(
-            fmt::format("data batch field count {} is less than required 
special field count {}",
-                        data_batch->num_fields(), 
SpecialFields::KEY_VALUE_SPECIAL_FIELD_COUNT));
-    }
-    if (!data_batch->field(0) || data_batch->field(0)->type_id() != 
arrow::Type::INT64) {
+    std::shared_ptr<arrow::Array> sequence_number =
+        data_batch->GetFieldByName(SpecialFields::SequenceNumber().Name());
+    if (!sequence_number || sequence_number->type_id() != arrow::Type::INT64) {
         return Status::Invalid("cannot cast SEQUENCE_NUMBER column to int64 
arrow array");
     }
     sequence_number_array_ =
-        
checked_pointer_cast<arrow::NumericArray<arrow::Int64Type>>(data_batch->field(0));
-    if (!data_batch->field(1) || data_batch->field(1)->type_id() != 
arrow::Type::INT8) {
+        
checked_pointer_cast<arrow::NumericArray<arrow::Int64Type>>(sequence_number);
+    if (sequence_number_array_->null_count() != 0) {
+        return Status::Invalid("SEQUENCE_NUMBER column contains null");
+    }
+    std::shared_ptr<arrow::Array> row_kind =
+        data_batch->GetFieldByName(SpecialFields::ValueKind().Name());
+    if (!row_kind || row_kind->type_id() != arrow::Type::INT8) {
         return Status::Invalid("cannot cast VALUE_KIND column to int8 arrow 
array");
     }
-    row_kind_array_ =
-        
checked_pointer_cast<arrow::NumericArray<arrow::Int8Type>>(data_batch->field(1));
+    row_kind_array_ = 
checked_pointer_cast<arrow::NumericArray<arrow::Int8Type>>(row_kind);
+    if (row_kind_array_->null_count() != 0) {
+        return Status::Invalid("VALUE_KIND column contains null");
+    }
     arrow::ArrayVector key_fields;
     key_fields.reserve(key_schema_->num_fields());
     for (const auto& key_field : key_schema_->fields()) {
-        // skip special fields
-        key_fields.emplace_back(data_batch->GetFieldByName(key_field->name()));
+        std::shared_ptr<arrow::Array> field_array = 
data_batch->GetFieldByName(key_field->name());

Review Comment:
   Done



##########
src/paimon/core/realtime/realtime_primary_key_reader.cpp:
##########
@@ -18,493 +18,74 @@
 
 #include "paimon/core/realtime/realtime_primary_key_reader.h"
 
-#include <cstdint>
 #include <memory>
-#include <optional>
-#include <unordered_map>
 #include <utility>
 #include <vector>
 
-#include "arrow/array/array_base.h"
-#include "arrow/array/array_primitive.h"
-#include "arrow/c/bridge.h"
 #include "arrow/type.h"
-#include "fmt/format.h"
-#include "paimon/common/data/columnar/columnar_batch_context.h"
-#include "paimon/common/data/columnar/columnar_row_ref.h"
 #include "paimon/common/table/special_fields.h"
 #include "paimon/common/types/data_field.h"
-#include "paimon/common/types/row_kind.h"
-#include "paimon/common/utils/arrow/arrow_utils.h"
-#include "paimon/common/utils/arrow/status_utils.h"
-#include "paimon/common/utils/checked_cast.h"
-#include "paimon/common/utils/scope_guard.h"
-#include "paimon/core/utils/nested_projection_utils.h"
-#include "paimon/macros.h"
-#include "paimon/reader/batch_reader.h"
+#include "paimon/core/io/key_value_data_file_record_reader.h"
+#include "paimon/core/key_value.h"
+#include "paimon/core/realtime/realtime_store_read_pipeline.h"
 #include "paimon/status.h"
-#include "paimon/utils/roaring_bitmap64.h"
 
 namespace paimon {
-
 namespace {
 
-template <typename Reader>
-void CloseReaders(const std::vector<std::unique_ptr<Reader>>& readers) {
-    for (const std::unique_ptr<Reader>& reader : readers) {
-        if (reader) {
-            reader->Close();
-        }
-    }
-}
-
-class RealtimeOffsetCoverage {
- public:
-    static Result<std::shared_ptr<RealtimeOffsetCoverage>> Create(const 
OffsetRange& offsets,
-                                                                  size_t 
reader_count,
-                                                                  bool 
allow_committed_prefix) {
-        if (offsets.begin < 0 || offsets.end < offsets.begin) {
-            return Status::Invalid("PK real-time store returned an invalid 
offset range");
-        }
-        return std::shared_ptr<RealtimeOffsetCoverage>(
-            new RealtimeOffsetCoverage(offsets, reader_count, 
allow_committed_prefix));
-    }
-
-    Status Add(const arrow::Int64Array& offsets) {
-        for (int64_t row = 0; row < offsets.length(); ++row) {
-            const int64_t offset = offsets.Value(row);
-            if (allow_committed_prefix_ && offset < 0) {
-                return Status::Invalid("PK real-time store reader offset must 
be non-negative");
-            }
-            if (allow_committed_prefix_ && offset < offsets_.begin) {
-                continue;
-            }
-            if (offset < offsets_.begin || offset >= offsets_.end) {
-                return Status::Invalid(
-                    allow_committed_prefix_
-                        ? "PK real-time store query reader offset is outside 
the visible range"
-                        : "PK real-time store commit reader offset is outside 
the sealed range");
-            }
-            if (!seen_offsets_.CheckedAdd(offset)) {
-                return CoverageError();
-            }
-        }
-        return Status::OK();
-    }
-
-    Status FinishReader() {
-        ++finished_reader_count_;
-        if (finished_reader_count_ == reader_count_ &&
-            seen_offsets_.Cardinality() != offsets_.Count()) {
-            return CoverageError();
-        }
-        return Status::OK();
-    }
-
- private:
-    RealtimeOffsetCoverage(const OffsetRange& offsets, size_t reader_count,
-                           bool allow_committed_prefix)
-        : offsets_(offsets),
-          reader_count_(reader_count),
-          allow_committed_prefix_(allow_committed_prefix) {}
-
-    Status CoverageError() const {
-        return Status::Invalid(
-            allow_committed_prefix_
-                ? "PK real-time store query readers did not cover the visible 
range"
-                : "PK real-time store commit readers did not cover the sealed 
range");
-    }
-
-    OffsetRange offsets_;
-    size_t reader_count_;
-    bool allow_committed_prefix_;
-    RoaringBitmap64 seen_offsets_;
-    size_t finished_reader_count_ = 0;
-};
-
-Status CheckTransportField(const std::shared_ptr<arrow::Schema>& schema, 
int32_t field_idx,
-                           const DataField& expected_field) {
-    if (schema->num_fields() <= field_idx) {
-        return Status::Invalid(
-            fmt::format("realtime primary-key transport schema is missing 
field {} at index {}",
-                        expected_field.Name(), field_idx));
-    }
-    const std::shared_ptr<arrow::Field>& field = schema->field(field_idx);
-    PAIMON_ASSIGN_OR_RAISE(int32_t field_id, 
NestedProjectionUtils::GetPaimonFieldId(field));
-    if (field->name() != expected_field.Name() || 
!field->type()->Equals(*expected_field.Type()) ||
-        field->nullable() || field_id != expected_field.Id()) {
-        return Status::Invalid(fmt::format(
-            "realtime primary-key transport schema field {} must be non-null 
{}:{} with field id "
-            "{}, got {}:{} nullable={} field id {}",
-            field_idx, expected_field.Name(), 
expected_field.Type()->ToString(),
-            expected_field.Id(), field->name(), field->type()->ToString(), 
field->nullable(),
-            field_id));
-    }
-    return Status::OK();
-}
-
-Result<std::vector<int32_t>> ResolveFieldIndexes(
-    const std::shared_ptr<arrow::Schema>& transport_schema,
-    const std::unordered_map<int32_t, int32_t>& field_indexes,
-    const std::shared_ptr<arrow::Schema>& row_schema) {
-    std::vector<int32_t> result;
-    result.reserve(row_schema->num_fields());
-    for (const std::shared_ptr<arrow::Field>& row_field : 
row_schema->fields()) {
-        PAIMON_ASSIGN_OR_RAISE(int32_t field_id,
-                               
NestedProjectionUtils::GetPaimonFieldId(row_field));
-        auto field_index = field_indexes.find(field_id);
-        if (field_index == field_indexes.end()) {
-            return Status::Invalid(fmt::format(
-                "cannot find field id {} in realtime primary-key transport 
schema", field_id));
-        }
-        const std::shared_ptr<arrow::Field>& transport_field =
-            transport_schema->field(field_index->second);
-        if (!transport_field->type()->Equals(row_field->type())) {
-            return Status::Invalid(fmt::format(
-                "realtime primary-key transport field id {} type {} does not 
match row type {}",
-                field_id, transport_field->type()->ToString(), 
row_field->type()->ToString()));
+Result<std::vector<std::unique_ptr<KeyValueRecordReader>>> 
CreateKeyValueReaders(
+    std::vector<std::unique_ptr<BatchReader>>&& readers,
+    const std::shared_ptr<arrow::Schema>& key_schema,
+    const std::shared_ptr<arrow::Schema>& value_schema,
+    const std::shared_ptr<MemoryPool>& memory_pool) {
+    std::vector<std::unique_ptr<KeyValueRecordReader>> result;
+    result.reserve(readers.size());
+    for (std::unique_ptr<BatchReader>& reader : readers) {

Review Comment:
   Done



-- 
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