| 1647 | } |
| 1648 | |
| 1649 | Status TmpFileGroup::ReadAsync(TmpWriteHandle* handle, MemRange buffer) { |
| 1650 | DCHECK(handle->write_range_ != nullptr); |
| 1651 | DCHECK(!handle->is_cancelled_); |
| 1652 | DCHECK_EQ(buffer.len(), handle->data_len()); |
| 1653 | Status status; |
| 1654 | VLOG(3) << "ReadAsync " << handle->TmpFilePath() << " " |
| 1655 | << handle->write_range_->offset() << " " << handle->on_disk_len(); |
| 1656 | // Don't grab 'write_state_lock_' in this method - it is not necessary because we |
| 1657 | // don't touch any members that it protects and could block other threads for the |
| 1658 | // duration of the synchronous read. |
| 1659 | DCHECK(!handle->write_in_flight_); |
| 1660 | DCHECK(handle->read_range_ == nullptr); |
| 1661 | DCHECK(handle->write_range_ != nullptr); |
| 1662 | |
| 1663 | MemRange read_buffer = buffer; |
| 1664 | if (handle->is_compressed()) { |
| 1665 | int64_t compressed_len = handle->compressed_len_; |
| 1666 | if (!handle->compressed_.TryAllocate(compressed_len)) { |
| 1667 | return tmp_file_mgr_->compressed_buffer_tracker()->MemLimitExceeded( |
| 1668 | nullptr, "Failed to decompress spilled data", compressed_len); |
| 1669 | } |
| 1670 | DCHECK_EQ(compressed_len, handle->write_range_->len()); |
| 1671 | read_buffer = MemRange(handle->compressed_.buffer(), compressed_len); |
| 1672 | } |
| 1673 | |
| 1674 | // Don't grab handle->write_state_lock_, it is safe to touch all of handle's state |
| 1675 | // since the write is not in flight. |
| 1676 | handle->read_range_ = scan_range_pool_.Add(new ScanRange); |
| 1677 | int64_t offset = handle->write_range_->offset(); |
| 1678 | if (handle->file_ != nullptr && !handle->file_->is_local()) { |
| 1679 | TmpFileRemote* tmp_file = static_cast<TmpFileRemote*>(handle->file_); |
| 1680 | DiskFile* local_read_buffer_file = tmp_file->GetReadBufferFile(offset); |
| 1681 | DiskFile* remote_file = tmp_file->DiskFile(); |
| 1682 | // Reset the read_range, use the remote filesystem's disk id. |
| 1683 | handle->read_range_->Reset( |
| 1684 | ScanRange::FileInfo{ |
| 1685 | remote_file->path().c_str(), tmp_file->hdfs_conn_, tmp_file->mtime_}, |
| 1686 | handle->write_range_->len(), offset, tmp_file->disk_id(), false, |
| 1687 | BufferOpts::ReadInto( |
| 1688 | read_buffer.data(), read_buffer.len(), BufferOpts::NO_CACHING), |
| 1689 | nullptr, remote_file, local_read_buffer_file); |
| 1690 | } else { |
| 1691 | // Read from local. |
| 1692 | handle->read_range_->Reset( |
| 1693 | ScanRange::FileInfo{handle->write_range_->file()}, |
| 1694 | handle->write_range_->len(), offset, handle->write_range_->disk_id(), false, |
| 1695 | BufferOpts::ReadInto( |
| 1696 | read_buffer.data(), read_buffer.len(), BufferOpts::NO_CACHING)); |
| 1697 | } |
| 1698 | |
| 1699 | read_counter_->Add(1); |
| 1700 | bytes_read_counter_->Add(read_buffer.len()); |
| 1701 | |
| 1702 | bool needs_buffers; |
| 1703 | RETURN_IF_ERROR(io_ctx_->StartScanRange(handle->read_range_, &needs_buffers)); |
| 1704 | DCHECK(!needs_buffers) << "Already provided a buffer"; |
| 1705 | return Status::OK(); |
| 1706 | } |
no test coverage detected