wangyong9999 commented on code in PR #314:
URL: https://github.com/apache/paimon-cpp/pull/314#discussion_r3985480871
##########
src/paimon/format/parquet/parquet_input_stream.h:
##########
@@ -78,21 +80,22 @@ class ParquetInputStream : public ArrowInputStreamAdapter {
using ArrowInputStreamAdapter::ReadAt;
arrow::Result<int64_t> ReadAt(int64_t position, int64_t nbytes, void* out)
override {
- if (!cache_ || file_uri_.empty() || nbytes <= 0 ||
- nbytes > std::numeric_limits<int32_t>::max()) {
+ if (!cache_ || file_uri_.empty() || position < 0 || nbytes <= 0 ||
position > file_size_ ||
+ nbytes > file_size_ - position || nbytes >
std::numeric_limits<int32_t>::max()) {
return ArrowInputStreamAdapter::ReadAt(position, nbytes, out);
}
+ bool is_index = false;
auto range = index_ranges_.upper_bound(position);
- if (range == index_ranges_.begin()) {
- return ArrowInputStreamAdapter::ReadAt(position, nbytes, out);
+ if (range != index_ranges_.begin()) {
+ --range;
+ const int64_t offset = position - range->first;
+ is_index = offset <= range->second && nbytes <= range->second -
offset;
}
- --range;
- if (position - range->first > range->second ||
- nbytes > range->second - (position - range->first)) {
+ if (!is_index && !cache_data_) {
return ArrowInputStreamAdapter::ReadAt(position, nbytes, out);
}
auto key = CacheKey::ForKind(file_uri_, position,
static_cast<int32_t>(nbytes),
Review Comment:
The key here is whatever range Arrow happened to request: with pre-buffer on
that is the coalesced range for this projection and page selection, so the same
bytes read under a different projection or selection get a new entry instead of
a hit, and with pre-buffer off a whole column chunk becomes one entry. There is
no cap on entry size, so a plain `LruCache` shared with the footer/index
entries gets flushed by one large chunk read.
Main just merged #272, which puts a `FileBlockCache` one layer below this
(fixed-size blocks, one fetch per block, concurrent readers wait on the same
fetch); this PR is based before it. The only thing that cache lacks is reuse
across reader lifetimes. Backing its blocks with the caller-provided `Cache`
(uri + block index) would give that with bounded entries and hits on
overlapping reads, instead of a second cache over the same bytes keyed by
request shape. At minimum this needs a rebase onto #272 and a note on how the
two layers interact.
##########
src/paimon/format/parquet/parquet_input_stream.h:
##########
@@ -103,24 +106,89 @@ class ParquetInputStream : public ArrowInputStreamAdapter
{
int64_t size,
ArrowInputStreamAdapter::ReadAt(position, nbytes,
segment.MutableData()));
if (size != nbytes) {
- return Status::IOError("Short read of Parquet page index");
+ return Status::IOError("Short read of Parquet cached
range");
}
return std::make_shared<CacheValue>(segment, CacheCallback());
});
if (!value.ok()) {
return ToArrowStatus(value.status());
}
if (!value.value() || value.value()->GetSegment().Size() != nbytes) {
- return arrow::Status::IOError("Invalid Parquet page-index cache
value");
+ return arrow::Status::IOError("Invalid Parquet cached range
value");
}
std::memcpy(out, value.value()->GetSegment().Data(), nbytes);
return nbytes;
}
+ arrow::Future<std::shared_ptr<arrow::Buffer>> ReadAsync(const
arrow::io::IOContext& io_context,
+ int64_t position,
+ int64_t nbytes)
override {
+ if (!cache_data_ || !cache_ || file_uri_.empty() || position < 0 ||
nbytes <= 0 ||
Review Comment:
This guard and the `is_index` lookup below are a copy of the ones in
`ReadAt`. One helper returning the cache key (or nullopt when the read is not
cacheable) would serve both.
--
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]