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 highlighting the potential overhead on the I/O hot path. I’ve
refined the implementation in two aspects:
1. The latency histograms have been replaced with lightweight count and
sum-us counters. Their atomic updates use std::memory_order_relaxed, minimizing
the runtime cost of metrics collection.
2. A new prefetch.io-metrics.enabled option has been introduced and is
disabled by default. When it is disabled, MetricsInputStream is not
instantiated at all, ensuring that the underlying I/O path remains free from
instrumentation overhead.
--
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]