ClickHouse/dbms/include/DB/IO/WriteBufferFromFileDescriptor.h

130 lines
3.1 KiB
C
Raw Normal View History

2011-10-24 12:10:59 +00:00
#pragma once
#include <unistd.h>
#include <errno.h>
2015-10-05 01:35:28 +00:00
#include <DB/Common/Exception.h>
#include <DB/Common/ProfileEvents.h>
#include <DB/Common/CurrentMetrics.h>
2011-10-24 12:10:59 +00:00
#include <DB/IO/WriteBufferFromFileBase.h>
2011-10-24 12:10:59 +00:00
#include <DB/IO/WriteBuffer.h>
2013-06-21 20:34:19 +00:00
#include <DB/IO/WriteHelpers.h>
2011-10-24 12:10:59 +00:00
#include <DB/IO/BufferWithOwnMemory.h>
namespace DB
{
namespace ErrorCodes
{
extern const int CANNOT_WRITE_TO_FILE_DESCRIPTOR;
extern const int CANNOT_FSYNC;
extern const int CANNOT_SEEK_THROUGH_FILE;
extern const int CANNOT_TRUNCATE_FILE;
}
2011-10-24 12:10:59 +00:00
/** Работает с готовым файловым дескриптором. Не открывает и не закрывает файл.
*/
class WriteBufferFromFileDescriptor : public WriteBufferFromFileBase
2011-10-24 12:10:59 +00:00
{
protected:
int fd;
2016-03-07 04:31:10 +00:00
void nextImpl() override
2011-10-24 12:10:59 +00:00
{
if (!offset())
return;
2011-12-26 07:07:30 +00:00
size_t bytes_written = 0;
while (bytes_written != offset())
{
ProfileEvents::increment(ProfileEvents::WriteBufferFromFileDescriptorWrite);
ssize_t res = 0;
{
CurrentMetrics::Increment metric_increment{CurrentMetrics::Write};
res = ::write(fd, working_buffer.begin() + bytes_written, offset() - bytes_written);
}
2011-12-26 07:07:30 +00:00
if ((-1 == res || 0 == res) && errno != EINTR)
throwFromErrno("Cannot write to file " + getFileName(), ErrorCodes::CANNOT_WRITE_TO_FILE_DESCRIPTOR);
if (res > 0)
bytes_written += res;
}
ProfileEvents::increment(ProfileEvents::WriteBufferFromFileDescriptorWriteBytes, bytes_written);
2011-10-24 12:10:59 +00:00
}
/// Имя или описание файла
virtual std::string getFileName() const override
2011-10-24 12:10:59 +00:00
{
2013-06-21 20:34:19 +00:00
return "(fd = " + toString(fd) + ")";
2011-10-24 12:10:59 +00:00
}
public:
2014-04-08 07:31:51 +00:00
WriteBufferFromFileDescriptor(int fd_ = -1, size_t buf_size = DBMS_DEFAULT_BUFFER_SIZE, char * existing_memory = nullptr, size_t alignment = 0)
: WriteBufferFromFileBase(buf_size, existing_memory, alignment), fd(fd_) {}
2011-10-24 12:10:59 +00:00
/** Можно вызывать для инициализации, если нужный fd не был передан в конструктор.
* Менять fd во время работы нельзя.
*/
void setFD(int fd_)
{
fd = fd_;
}
~WriteBufferFromFileDescriptor()
2011-10-24 12:10:59 +00:00
{
2013-11-18 17:17:45 +00:00
try
{
2015-12-13 08:51:28 +00:00
if (fd >= 0)
next();
2013-11-18 17:17:45 +00:00
}
catch (...)
{
2013-11-18 19:18:03 +00:00
tryLogCurrentException(__PRETTY_FUNCTION__);
2013-11-18 17:17:45 +00:00
}
2011-10-24 12:10:59 +00:00
}
int getFD() const override
{
return fd;
}
off_t getPositionInFile() override
{
return seek(0, SEEK_CUR);
}
void sync() override
2013-09-15 01:10:16 +00:00
{
/// Если в буфере ещё остались данные - запишем их.
next();
/// Попросим ОС сбросить данные на диск.
int res = fsync(fd);
if (-1 == res)
throwFromErrno("Cannot fsync " + getFileName(), ErrorCodes::CANNOT_FSYNC);
}
private:
off_t doSeek(off_t offset, int whence) override
{
off_t res = lseek(fd, offset, whence);
if (-1 == res)
throwFromErrno("Cannot seek through file " + getFileName(), ErrorCodes::CANNOT_SEEK_THROUGH_FILE);
return res;
}
void doTruncate(off_t length) override
{
int res = ftruncate(fd, length);
if (-1 == res)
throwFromErrno("Cannot truncate file " + getFileName(), ErrorCodes::CANNOT_TRUNCATE_FILE);
}
2011-10-24 12:10:59 +00:00
};
}