gripleaf commented on code in PR #211:
URL: https://github.com/apache/paimon-cpp/pull/211#discussion_r3868419689


##########
src/paimon/common/reader/prefetch_file_batch_reader_impl.cpp:
##########
@@ -42,8 +42,135 @@ class Schema;
 
 namespace paimon {
 
+struct PrefetchMetricsState {
+    std::atomic<uint64_t> read_ranges_total{0};
+    std::atomic<uint64_t> read_ranges_after_bitmap{0};
+    std::atomic<uint64_t> seek_count{0};
+    std::atomic<uint64_t> produced_batches{0};
+    std::atomic<uint64_t> consumed_batches{0};
+    std::atomic<uint64_t> discarded_batches{0};
+    std::atomic<uint64_t> errors{0};
+    std::atomic<uint64_t> adaptive_disabled_count{0};
+    std::atomic<uint64_t> queue_full_count{0};
+    std::atomic<uint64_t> queue_depth{0};
+    std::atomic<uint64_t> queue_depth_max{0};
+    std::atomic<bool> enabled{false};
+    std::shared_ptr<MetricsImpl> histograms = std::make_shared<MetricsImpl>();
+};
+
+struct PrefetchIoMetricsState {
+    std::atomic<uint64_t> read_requests{0};
+    std::atomic<uint64_t> read_requested_bytes{0};
+    std::atomic<uint64_t> read_physical_bytes{0};
+    std::atomic<uint64_t> read_failed{0};
+    std::atomic<uint64_t> async_requests{0};
+    std::atomic<uint64_t> async_requested_bytes{0};
+    std::atomic<uint64_t> async_physical_bytes{0};
+    std::atomic<uint64_t> async_completed{0};
+    std::atomic<uint64_t> async_failed{0};
+    std::atomic<uint64_t> async_pending{0};
+    std::shared_ptr<MetricsImpl> histograms = std::make_shared<MetricsImpl>();
+};
+
 namespace {
 
+uint64_t ElapsedMicros(const std::chrono::steady_clock::time_point& start) {
+    return 
static_cast<uint64_t>(std::chrono::duration_cast<std::chrono::microseconds>(
+                                     std::chrono::steady_clock::now() - start)
+                                     .count());
+}
+
+void UpdateMax(std::atomic<uint64_t>* target, uint64_t value) {
+    uint64_t current = target->load();
+    while (current < value && !target->compare_exchange_weak(current, value)) {
+    }
+}
+
+class MetricsInputStream : public InputStream {
+ public:
+    MetricsInputStream(const std::shared_ptr<InputStream>& stream,
+                       const std::shared_ptr<PrefetchIoMetricsState>& metrics)
+        : stream_(stream), metrics_(metrics) {}
+
+    MetricsInputStream(std::unique_ptr<InputStream>&& stream,
+                       const std::shared_ptr<PrefetchIoMetricsState>& metrics)
+        : stream_(std::move(stream)), metrics_(metrics) {}
+
+    Status Seek(int64_t offset, SeekOrigin origin) override {
+        return stream_->Seek(offset, origin);
+    }
+
+    Result<int64_t> GetPos() const override {
+        return stream_->GetPos();
+    }
+
+    Result<int64_t> Read(char* buffer, int64_t size) override {
+        return RecordRead([&]() { return stream_->Read(buffer, size); }, size);
+    }
+
+    Result<int64_t> Read(char* buffer, int64_t size, int64_t offset) override {
+        return RecordRead([&]() { return stream_->Read(buffer, size, offset); 
}, size);
+    }
+
+    void ReadAsync(char* buffer, int64_t size, int64_t offset,
+                   std::function<void(Status)>&& callback) override {
+        metrics_->async_requests.fetch_add(1);
+        metrics_->async_requested_bytes.fetch_add(
+            static_cast<uint64_t>(std::max<int64_t>(0, size)));
+        metrics_->async_pending.fetch_add(1);
+        std::shared_ptr<PrefetchIoMetricsState> metrics = metrics_;
+        const auto start = std::chrono::steady_clock::now();
+        stream_->ReadAsync(
+            buffer, size, offset,
+            [metrics, size, start, callback = std::move(callback)](Status 
status) mutable {
+                metrics->async_pending.fetch_sub(1);
+                if (status.ok()) {
+                    metrics->async_completed.fetch_add(1);
+                    metrics->async_physical_bytes.fetch_add(
+                        static_cast<uint64_t>(std::max<int64_t>(0, size)));
+                } else {
+                    metrics->async_failed.fetch_add(1);
+                }
+                
metrics->histograms->ObserveHistogram(PrefetchIoMetrics::ASYNC_LATENCY_US,
+                                                      ElapsedMicros(start));
+                callback(status);
+            });
+    }
+
+    Status Close() override {
+        return stream_->Close();
+    }
+
+    Result<std::string> GetUri() const override {
+        return stream_->GetUri();
+    }
+
+    Result<int64_t> Length() const override {
+        return stream_->Length();
+    }
+
+ private:
+    template <typename ReadFunction>
+    Result<int64_t> RecordRead(ReadFunction&& read, int64_t size) {
+        metrics_->read_requests.fetch_add(1);

Review Comment:
   Thanks for pointing out the potential hot-path overhead. I’ve addressed it 
in two ways:
   
     1. Replaced the latency histograms with count and sum-us counters, using 
std::memory_order_relaxed for the atomic operations to reduce metrics overhead.
     2. Added the prefetch.io-metrics.enabled option, which is disabled by 
default. When disabled, MetricsInputStream is not constructed, physically 
avoiding any performance overhead from the per-I/O metrics instrumentation.



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