diff options
| author | Vasily Gerasimov <[email protected]> | 2025-02-21 09:48:22 +0200 |
|---|---|---|
| committer | GitHub <[email protected]> | 2025-02-21 09:48:22 +0200 |
| commit | 499b85ea8bc2a2d6f99427c81928535dbba83e1b (patch) | |
| tree | 60ca17406eb3796cec0074633c0b687bba9f8a9c | |
| parent | 68ccd290e81c8e8be4673a58dc03439fdadacd2a (diff) | |
Refactor S3Buffer for export so that it can use different combinations of checksums, compression and encryption (#14862)
| -rw-r--r-- | ydb/core/tx/datashard/export_s3.h | 19 | ||||
| -rw-r--r-- | ydb/core/tx/datashard/export_s3_buffer.cpp | 244 | ||||
| -rw-r--r-- | ydb/core/tx/datashard/export_s3_buffer.h | 94 | ||||
| -rw-r--r-- | ydb/core/tx/datashard/export_s3_buffer_raw.h | 55 | ||||
| -rw-r--r-- | ydb/core/tx/datashard/export_s3_buffer_zstd.cpp | 138 | ||||
| -rw-r--r-- | ydb/core/tx/datashard/ya.make | 1 |
6 files changed, 325 insertions, 226 deletions
diff --git a/ydb/core/tx/datashard/export_s3.h b/ydb/core/tx/datashard/export_s3.h index cf21cd2e90d..ddb843eb749 100644 --- a/ydb/core/tx/datashard/export_s3.h +++ b/ydb/core/tx/datashard/export_s3.h @@ -29,15 +29,28 @@ public: const ui64 maxBytes = scanSettings.GetBytesBatchSize(); const ui64 minBytes = Task.GetS3Settings().GetLimits().GetMinWriteBatchSize(); + TS3ExportBufferSettings bufferSettings; + bufferSettings + .WithColumns(Columns) + .WithMaxRows(maxRows) + .WithMaxBytes(maxBytes); + if (Task.GetEnableChecksums()) { + bufferSettings.WithChecksum(TS3ExportBufferSettings::Sha256Checksum()); + } + switch (CodecFromTask(Task)) { case ECompressionCodec::None: - return CreateS3ExportBufferRaw(Columns, maxRows, maxBytes, Task.GetEnableChecksums()); + break; case ECompressionCodec::Zstd: - return CreateS3ExportBufferZstd(Task.GetCompression().GetLevel(), Columns, maxRows, - maxBytes, minBytes, Task.GetEnableChecksums()); + bufferSettings + .WithMinBytes(minBytes) + .WithCompression(TS3ExportBufferSettings::ZstdCompression(Task.GetCompression().GetLevel())); + break; case ECompressionCodec::Invalid: Y_ABORT("unreachable"); } + + return CreateS3ExportBuffer(std::move(bufferSettings)); } void Shutdown() const override {} diff --git a/ydb/core/tx/datashard/export_s3_buffer.cpp b/ydb/core/tx/datashard/export_s3_buffer.cpp index 2af4b591af0..6a97fa4bcd9 100644 --- a/ydb/core/tx/datashard/export_s3_buffer.cpp +++ b/ydb/core/tx/datashard/export_s3_buffer.cpp @@ -1,8 +1,9 @@ #ifndef KIKIMR_DISABLE_S3_OPS -#include "export_s3_buffer_raw.h" +#include "export_s3_buffer.h" #include "type_serialization.h" +#include <ydb/core/backup/common/checksum.h> #include <ydb/core/tablet_flat/flat_row_state.h> #include <yql/essentials/types/binary_json/read.h> #include <ydb/public/lib/scheme_types/scheme_type_id.h> @@ -10,22 +11,128 @@ #include <library/cpp/string_utils/quote/quote.h> #include <util/datetime/base.h> +#include <util/generic/buffer.h> #include <util/stream/buffer.h> -namespace NKikimr { -namespace NDataShard { +#include <contrib/libs/zstd/include/zstd.h> -TS3BufferRaw::TS3BufferRaw(const TTagToColumn& columns, ui64 rowsLimit, ui64 bytesLimit, bool enableChecksums) - : Columns(columns) - , RowsLimit(rowsLimit) - , BytesLimit(bytesLimit) - , Rows(0) - , BytesRead(0) - , Checksum(enableChecksums ? NBackup::CreateChecksum() : nullptr) + +namespace NKikimr::NDataShard { + +namespace { + +struct DestroyZCtx { + static void Destroy(::ZSTD_CCtx* p) noexcept { + ZSTD_freeCCtx(p); + } +}; + +class TZStdCompressionProcessor { +public: + using TPtr = THolder<TZStdCompressionProcessor>; + + explicit TZStdCompressionProcessor(const TS3ExportBufferSettings::TCompressionSettings& settings); + + TString GetError() const { + return ZSTD_getErrorName(ErrorCode); + } + + bool AddData(TStringBuf data); + + TMaybe<TBuffer> Flush(bool prepare); + +private: + enum ECompressionResult { + CONTINUE, + DONE, + ERROR, + }; + + ECompressionResult Compress(ZSTD_inBuffer* input, ZSTD_EndDirective endOp); + void Reset(); + +private: + const int CompressionLevel; + THolder<::ZSTD_CCtx, DestroyZCtx> Context; + size_t ErrorCode = 0; + TBuffer Buffer; + ui64 BytesAdded = 0; +}; + +class TS3Buffer: public NExportScan::IBuffer { + using TTagToColumn = IExport::TTableColumns; + using TTagToIndex = THashMap<ui32, ui32>; // index in IScan::TRow + +public: + explicit TS3Buffer(TS3ExportBufferSettings&& settings); + + void ColumnsOrder(const TVector<ui32>& tags) override; + bool Collect(const NTable::IScan::TRow& row) override; + IEventBase* PrepareEvent(bool last, NExportScan::IBuffer::TStats& stats) override; + void Clear() override; + bool IsFilled() const override; + TString GetError() const override; + +private: + inline ui64 GetRowsLimit() const { return RowsLimit; } + inline ui64 GetBytesLimit() const { return MaxBytes; } + + bool Collect(const NTable::IScan::TRow& row, IOutputStream& out); + virtual TMaybe<TBuffer> Flush(bool prepare); + + static NBackup::IChecksum* CreateChecksum(const TMaybe<TS3ExportBufferSettings::TChecksumSettings>& settings); + static TZStdCompressionProcessor* CreateCompression(const TMaybe<TS3ExportBufferSettings::TCompressionSettings>& settings); + +private: + const TTagToColumn Columns; + const ui64 RowsLimit; + const ui64 MinBytes; + const ui64 MaxBytes; + + TTagToIndex Indices; + +protected: + ui64 Rows = 0; + ui64 BytesRead = 0; + TBuffer Buffer; + + NBackup::IChecksum::TPtr Checksum; + TZStdCompressionProcessor::TPtr Compression; + + TString ErrorString; +}; // TS3Buffer + +TS3Buffer::TS3Buffer(TS3ExportBufferSettings&& settings) + : Columns(std::move(settings.Columns)) + , RowsLimit(settings.MaxRows) + , MinBytes(settings.MinBytes) + , MaxBytes(settings.MaxBytes) + , Checksum(CreateChecksum(settings.ChecksumSettings)) + , Compression(CreateCompression(settings.CompressionSettings)) { } -void TS3BufferRaw::ColumnsOrder(const TVector<ui32>& tags) { +NBackup::IChecksum* TS3Buffer::CreateChecksum(const TMaybe<TS3ExportBufferSettings::TChecksumSettings>& settings) { + if (settings) { + switch (settings->ChecksumType) { + case TS3ExportBufferSettings::TChecksumSettings::EChecksumType::Sha256: + return NBackup::CreateChecksum(); + } + } + return nullptr; +} + +TZStdCompressionProcessor* TS3Buffer::CreateCompression(const TMaybe<TS3ExportBufferSettings::TCompressionSettings>& settings) { + if (settings) { + switch (settings->Algorithm) { + case TS3ExportBufferSettings::TCompressionSettings::EAlgorithm::Zstd: + return new TZStdCompressionProcessor(*settings); + } + } + return nullptr; +} + +void TS3Buffer::ColumnsOrder(const TVector<ui32>& tags) { Y_ABORT_UNLESS(tags.size() == Columns.size()); Indices.clear(); @@ -37,7 +144,7 @@ void TS3BufferRaw::ColumnsOrder(const TVector<ui32>& tags) { } } -bool TS3BufferRaw::Collect(const NTable::IScan::TRow& row, IOutputStream& out) { +bool TS3Buffer::Collect(const NTable::IScan::TRow& row, IOutputStream& out) { bool needsComma = false; for (const auto& [tag, column] : Columns) { auto it = Indices.find(tag); @@ -152,7 +259,7 @@ bool TS3BufferRaw::Collect(const NTable::IScan::TRow& row, IOutputStream& out) { return true; } -bool TS3BufferRaw::Collect(const NTable::IScan::TRow& row) { +bool TS3Buffer::Collect(const NTable::IScan::TRow& row) { TBufferOutput out(Buffer); ErrorString.clear(); @@ -161,14 +268,24 @@ bool TS3BufferRaw::Collect(const NTable::IScan::TRow& row) { return false; } + TStringBuf data(Buffer.Data(), Buffer.Size()); + data = data.Tail(beforeSize); + + // Apply checksum if (Checksum) { - TStringBuf data(Buffer.Data(), Buffer.Size()); - Checksum->AddData(data.Tail(beforeSize)); + Checksum->AddData(data); + } + + // Compress + if (Compression && !Compression->AddData(data)) { + ErrorString = Compression->GetError(); + return false; } + return true; } -IEventBase* TS3BufferRaw::PrepareEvent(bool last, NExportScan::IBuffer::TStats& stats) { +IEventBase* TS3Buffer::PrepareEvent(bool last, NExportScan::IBuffer::TStats& stats) { stats.Rows = Rows; stats.BytesRead = BytesRead; @@ -186,31 +303,108 @@ IEventBase* TS3BufferRaw::PrepareEvent(bool last, NExportScan::IBuffer::TStats& } } -void TS3BufferRaw::Clear() { +void TS3Buffer::Clear() { Y_ABORT_UNLESS(Flush(false)); } -bool TS3BufferRaw::IsFilled() const { +bool TS3Buffer::IsFilled() const { + if (Buffer.Size() < MinBytes) { + return false; + } + return Rows >= GetRowsLimit() || Buffer.Size() >= GetBytesLimit(); } -TString TS3BufferRaw::GetError() const { +TString TS3Buffer::GetError() const { return ErrorString; } -TMaybe<TBuffer> TS3BufferRaw::Flush(bool) { +TMaybe<TBuffer> TS3Buffer::Flush(bool prepare) { Rows = 0; BytesRead = 0; + + // Compression finishes compression frame during Flush + // so that last table row borders equal to compression frame borders. + // This full finished block must then be encrypted so that encryption frame + // has the same borders. + // It allows to import data in batches and save its state during import. + + if (Compression) { + TMaybe<TBuffer> compressedBuffer = Compression->Flush(prepare); + if (!compressedBuffer) { + return Nothing(); + } + + Buffer = std::move(*compressedBuffer); + } + return std::exchange(Buffer, TBuffer()); } -NExportScan::IBuffer* CreateS3ExportBufferRaw( - const IExport::TTableColumns& columns, ui64 rowsLimit, ui64 bytesLimit, bool enableChecksums) +TZStdCompressionProcessor::TZStdCompressionProcessor(const TS3ExportBufferSettings::TCompressionSettings& settings) + : CompressionLevel(settings.CompressionLevel) + , Context(ZSTD_createCCtx()) { - return new TS3BufferRaw(columns, rowsLimit, bytesLimit, enableChecksums); } -} // NDataShard -} // NKikimr +bool TZStdCompressionProcessor::AddData(TStringBuf data) { + BytesAdded += data.size(); + auto input = ZSTD_inBuffer{data.data(), data.size(), 0}; + while (input.pos < input.size) { + if (ERROR == Compress(&input, ZSTD_e_continue)) { + return false; + } + } + + return true; +} + +TMaybe<TBuffer> TZStdCompressionProcessor::Flush(bool prepare) { + if (prepare && BytesAdded) { + ECompressionResult res; + auto input = ZSTD_inBuffer{NULL, 0, 0}; + + do { + if (res = Compress(&input, ZSTD_e_end); res == ERROR) { + return Nothing(); + } + } while (res != DONE); + } + + Reset(); + return std::exchange(Buffer, TBuffer()); +} + +TZStdCompressionProcessor::ECompressionResult TZStdCompressionProcessor::Compress(ZSTD_inBuffer* input, ZSTD_EndDirective endOp) { + auto output = ZSTD_outBuffer{Buffer.Data(), Buffer.Capacity(), Buffer.Size()}; + auto res = ZSTD_compressStream2(Context.Get(), &output, input, endOp); + + if (ZSTD_isError(res)) { + ErrorCode = res; + return ERROR; + } + + if (res > 0) { + Buffer.Reserve(output.pos + res); + } + + Buffer.Proceed(output.pos); + return res ? CONTINUE : DONE; +} + +void TZStdCompressionProcessor::Reset() { + BytesAdded = 0; + ZSTD_CCtx_reset(Context.Get(), ZSTD_reset_session_only); + ZSTD_CCtx_refCDict(Context.Get(), NULL); + ZSTD_CCtx_setParameter(Context.Get(), ZSTD_c_compressionLevel, CompressionLevel); +} + +} // anonymous + +NExportScan::IBuffer* CreateS3ExportBuffer(TS3ExportBufferSettings&& settings) { + return new TS3Buffer(std::move(settings)); +} + +} // namespace NKikimr::NDataShard #endif // KIKIMR_DISABLE_S3_OPS diff --git a/ydb/core/tx/datashard/export_s3_buffer.h b/ydb/core/tx/datashard/export_s3_buffer.h index 130e0ba7c96..9469899a5ee 100644 --- a/ydb/core/tx/datashard/export_s3_buffer.h +++ b/ydb/core/tx/datashard/export_s3_buffer.h @@ -5,14 +5,100 @@ #include "export_iface.h" #include "export_scan.h" +#include <util/generic/maybe.h> + namespace NKikimr { namespace NDataShard { -NExportScan::IBuffer* CreateS3ExportBufferRaw( - const IExport::TTableColumns& columns, ui64 maxRows, ui64 maxBytes, bool enableChecksums); +struct TS3ExportBufferSettings { + // Subsettings + struct TChecksumSettings { + enum class EChecksumType { + Sha256, + }; + + // Builders + TChecksumSettings& WithChecksumType(EChecksumType t) { + ChecksumType = t; + return *this; + } + + // Fields + EChecksumType ChecksumType = EChecksumType::Sha256; + }; + + static TChecksumSettings Sha256Checksum() { + return TChecksumSettings().WithChecksumType(TChecksumSettings::EChecksumType::Sha256); + } + + struct TCompressionSettings { + enum class EAlgorithm { + Zstd, + }; + + // Builders + TCompressionSettings& WithAlgorithm(EAlgorithm algorithm) { + Algorithm = algorithm; + return *this; + } + + TCompressionSettings& WithCompressionLevel(int level) { + CompressionLevel = level; + return *this; + } + + // Fields + EAlgorithm Algorithm = EAlgorithm::Zstd; + int CompressionLevel = -1; + }; + + static TCompressionSettings ZstdCompression(int level) { + return TCompressionSettings().WithAlgorithm(TCompressionSettings::EAlgorithm::Zstd).WithCompressionLevel(level); + } + + // Builders + TS3ExportBufferSettings& WithColumns(IExport::TTableColumns columns) { + Columns = std::move(columns); + return *this; + } + + TS3ExportBufferSettings& WithMaxRows(ui64 maxRows) { + MaxRows = maxRows; + return *this; + } + + TS3ExportBufferSettings& WithMinBytes(ui64 minBytes) { + MinBytes = minBytes; + return *this; + } + + TS3ExportBufferSettings& WithMaxBytes(ui64 maxBytes) { + MaxBytes = maxBytes; + return *this; + } + + TS3ExportBufferSettings& WithChecksum(TChecksumSettings settings) { + ChecksumSettings.ConstructInPlace(std::move(settings)); + return *this; + } + + TS3ExportBufferSettings& WithCompression(TCompressionSettings settings) { + CompressionSettings.ConstructInPlace(std::move(settings)); + return *this; + } + + // Fields + IExport::TTableColumns Columns; + ui64 MaxRows = 0; + ui64 MinBytes = 0; + ui64 MaxBytes = 0; + + // Data processing + TMaybe<TChecksumSettings> ChecksumSettings; + TMaybe<TCompressionSettings> CompressionSettings; +}; -NExportScan::IBuffer* CreateS3ExportBufferZstd(int compressionLevel, - const IExport::TTableColumns& columns, ui64 maxRows, ui64 maxBytes, ui64 minBytes, bool enableChecksums); +NExportScan::IBuffer* CreateS3ExportBuffer(TS3ExportBufferSettings&& settings); } // NDataShard } // NKikimr diff --git a/ydb/core/tx/datashard/export_s3_buffer_raw.h b/ydb/core/tx/datashard/export_s3_buffer_raw.h deleted file mode 100644 index 6591242358d..00000000000 --- a/ydb/core/tx/datashard/export_s3_buffer_raw.h +++ /dev/null @@ -1,55 +0,0 @@ -#pragma once - -#ifndef KIKIMR_DISABLE_S3_OPS - -#include "export_s3_buffer.h" - -#include <ydb/core/backup/common/checksum.h> - -#include <util/generic/buffer.h> - -namespace NKikimr { -namespace NDataShard { - -class TS3BufferRaw: public NExportScan::IBuffer { - using TTagToColumn = IExport::TTableColumns; - using TTagToIndex = THashMap<ui32, ui32>; // index in IScan::TRow - -public: - explicit TS3BufferRaw(const TTagToColumn& columns, ui64 rowsLimit, ui64 bytesLimit, bool enableChecksums); - - void ColumnsOrder(const TVector<ui32>& tags) override; - bool Collect(const NTable::IScan::TRow& row) override; - IEventBase* PrepareEvent(bool last, NExportScan::IBuffer::TStats& stats) override; - void Clear() override; - bool IsFilled() const override; - TString GetError() const override; - -protected: - inline ui64 GetRowsLimit() const { return RowsLimit; } - inline ui64 GetBytesLimit() const { return BytesLimit; } - - bool Collect(const NTable::IScan::TRow& row, IOutputStream& out); - virtual TMaybe<TBuffer> Flush(bool prepare); - -private: - const TTagToColumn Columns; - const ui64 RowsLimit; - const ui64 BytesLimit; - - TTagToIndex Indices; - -protected: - ui64 Rows; - ui64 BytesRead; - TBuffer Buffer; - - NBackup::IChecksum::TPtr Checksum; - - TString ErrorString; -}; // TS3BufferRaw - -} // NDataShard -} // NKikimr - -#endif // KIKIMR_DISABLE_S3_OPS diff --git a/ydb/core/tx/datashard/export_s3_buffer_zstd.cpp b/ydb/core/tx/datashard/export_s3_buffer_zstd.cpp deleted file mode 100644 index d1312b39847..00000000000 --- a/ydb/core/tx/datashard/export_s3_buffer_zstd.cpp +++ /dev/null @@ -1,138 +0,0 @@ -#ifndef KIKIMR_DISABLE_S3_OPS - -#include "export_s3_buffer_raw.h" - -#include <contrib/libs/zstd/include/zstd.h> - -namespace { - - struct DestroyZCtx { - static void Destroy(::ZSTD_CCtx* p) noexcept { - ZSTD_freeCCtx(p); - } - }; - -} // anonymous - -namespace NKikimr { -namespace NDataShard { - -class TS3BufferZstd: public TS3BufferRaw { - enum ECompressionResult { - CONTINUE, - DONE, - ERROR, - }; - - ECompressionResult Compress(ZSTD_inBuffer* input, ZSTD_EndDirective endOp) { - auto output = ZSTD_outBuffer{Buffer.Data(), Buffer.Capacity(), Buffer.Size()}; - auto res = ZSTD_compressStream2(Context.Get(), &output, input, endOp); - - if (ZSTD_isError(res)) { - ErrorCode = res; - return ERROR; - } - - if (res > 0) { - Buffer.Reserve(output.pos + res); - } - - Buffer.Proceed(output.pos); - return res ? CONTINUE : DONE; - } - - void Reset() { - ZSTD_CCtx_reset(Context.Get(), ZSTD_reset_session_only); - ZSTD_CCtx_refCDict(Context.Get(), NULL); - ZSTD_CCtx_setParameter(Context.Get(), ZSTD_c_compressionLevel, CompressionLevel); - } - -public: - explicit TS3BufferZstd(int compressionLevel, - const IExport::TTableColumns& columns, ui64 maxRows, ui64 maxBytes, ui64 minBytes, - bool enableChecksums) - : TS3BufferRaw(columns, maxRows, maxBytes, enableChecksums) - , CompressionLevel(compressionLevel) - , MinBytes(minBytes) - , Context(ZSTD_createCCtx()) - , ErrorCode(0) - , BytesRaw(0) - { - Reset(); - } - - bool Collect(const NTable::IScan::TRow& row) override { - BufferRaw.clear(); - TStringOutput out(BufferRaw); - if (!TS3BufferRaw::Collect(row, out)) { - return false; - } - - if (Checksum) { - Checksum->AddData(BufferRaw); - } - BytesRaw += BufferRaw.size(); - - auto input = ZSTD_inBuffer{BufferRaw.data(), BufferRaw.size(), 0}; - while (input.pos < input.size) { - if (ERROR == Compress(&input, ZSTD_e_continue)) { - return false; - } - } - - return true; - } - - bool IsFilled() const override { - if (Buffer.Size() < MinBytes) { - return false; - } - - return Rows >= GetRowsLimit() || BytesRaw >= GetBytesLimit(); - } - - TString GetError() const override { - return ZSTD_getErrorName(ErrorCode); - } - -protected: - TMaybe<TBuffer> Flush(bool prepare) override { - if (prepare && BytesRaw) { - ECompressionResult res; - auto input = ZSTD_inBuffer{NULL, 0, 0}; - - do { - if (res = Compress(&input, ZSTD_e_end); res == ERROR) { - return Nothing(); - } - } while (res != DONE); - } - - Reset(); - - BytesRaw = 0; - return TS3BufferRaw::Flush(prepare); - } - -private: - const int CompressionLevel; - const ui64 MinBytes; - - THolder<::ZSTD_CCtx, DestroyZCtx> Context; - size_t ErrorCode; - ui64 BytesRaw; - TString BufferRaw; - -}; // TS3BufferZstd - -NExportScan::IBuffer* CreateS3ExportBufferZstd(int compressionLevel, - const IExport::TTableColumns& columns, ui64 maxRows, ui64 maxBytes, ui64 minBytes, - bool enableChecksums) -{ - return new TS3BufferZstd(compressionLevel, columns, maxRows, maxBytes, minBytes, enableChecksums); -} - -} // NDataShard -} // NKikimr - -#endif // KIKIMR_DISABLE_S3_OPS diff --git a/ydb/core/tx/datashard/ya.make b/ydb/core/tx/datashard/ya.make index 1211a12a3b6..82fee3c74e7 100644 --- a/ydb/core/tx/datashard/ya.make +++ b/ydb/core/tx/datashard/ya.make @@ -294,7 +294,6 @@ IF (OS_WINDOWS) ELSE() SRCS( export_s3_buffer.cpp - export_s3_buffer_zstd.cpp export_s3_uploader.cpp extstorage_usage_config.cpp import_s3.cpp |
