2021-07-26 00:34:36 +00:00
|
|
|
#pragma once
|
|
|
|
|
|
|
|
#include <IO/ReadBufferFromFileBase.h>
|
|
|
|
#include <IO/AsynchronousReader.h>
|
|
|
|
#include <Interpreters/Context.h>
|
|
|
|
|
|
|
|
#include <optional>
|
|
|
|
#include <unistd.h>
|
|
|
|
|
|
|
|
|
|
|
|
namespace DB
|
|
|
|
{
|
|
|
|
|
|
|
|
/** Use ready file descriptor. Does not open or close a file.
|
|
|
|
*/
|
|
|
|
class AsynchronousReadBufferFromFileDescriptor : public ReadBufferFromFileBase
|
|
|
|
{
|
|
|
|
protected:
|
|
|
|
AsynchronousReaderPtr reader;
|
2021-08-16 00:00:32 +00:00
|
|
|
Int32 priority;
|
2021-07-26 00:34:36 +00:00
|
|
|
|
|
|
|
Memory<> prefetch_buffer;
|
2021-08-04 00:07:04 +00:00
|
|
|
std::future<IAsynchronousReader::Result> prefetch_future;
|
2021-07-26 00:34:36 +00:00
|
|
|
|
|
|
|
const size_t required_alignment = 0; /// For O_DIRECT both file offsets and memory addresses have to be aligned.
|
|
|
|
size_t file_offset_of_buffer_end = 0; /// What offset in file corresponds to working_buffer.end().
|
|
|
|
int fd;
|
|
|
|
|
|
|
|
bool nextImpl() override;
|
|
|
|
|
|
|
|
/// Name or some description of file.
|
|
|
|
std::string getFileName() const override;
|
|
|
|
|
2021-07-27 23:47:28 +00:00
|
|
|
void finalize();
|
|
|
|
|
2021-07-26 00:34:36 +00:00
|
|
|
public:
|
|
|
|
AsynchronousReadBufferFromFileDescriptor(
|
2021-08-16 00:00:32 +00:00
|
|
|
AsynchronousReaderPtr reader_, Int32 priority_,
|
2021-07-26 00:34:36 +00:00
|
|
|
int fd_, size_t buf_size = DBMS_DEFAULT_BUFFER_SIZE, char * existing_memory = nullptr, size_t alignment = 0)
|
|
|
|
: ReadBufferFromFileBase(buf_size, existing_memory, alignment),
|
2021-08-16 00:00:32 +00:00
|
|
|
reader(std::move(reader_)), priority(priority_), prefetch_buffer(buf_size, alignment), required_alignment(alignment), fd(fd_)
|
2021-07-26 00:34:36 +00:00
|
|
|
{
|
|
|
|
}
|
|
|
|
|
|
|
|
~AsynchronousReadBufferFromFileDescriptor() override;
|
|
|
|
|
|
|
|
void prefetch() override;
|
|
|
|
|
|
|
|
int getFD() const
|
|
|
|
{
|
|
|
|
return fd;
|
|
|
|
}
|
|
|
|
|
|
|
|
off_t getPosition() override
|
|
|
|
{
|
|
|
|
return file_offset_of_buffer_end - (working_buffer.end() - pos);
|
|
|
|
}
|
|
|
|
|
|
|
|
/// If 'offset' is small enough to stay in buffer after seek, then true seek in file does not happen.
|
|
|
|
off_t seek(off_t off, int whence) override;
|
|
|
|
|
|
|
|
/// Seek to the beginning, discarding already read data if any. Useful to reread file that changes on every read.
|
|
|
|
void rewind();
|
|
|
|
};
|
|
|
|
|
|
|
|
}
|
|
|
|
|