1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
|
#include "export_iface.h"
#include "backup_restore_common.h"
#include "execution_unit_ctors.h"
#include "export_scan.h"
#include "export_s3.h"
namespace NKikimr {
namespace NDataShard {
class TBackupUnit : public TBackupRestoreUnitBase<TEvDataShard::TEvCancelBackup> {
using IBuffer = NExportScan::IBuffer;
protected:
bool IsRelevant(TActiveTransaction* tx) const override {
return tx->GetSchemeTx().HasBackup();
}
bool IsWaiting(TOperation::TPtr op) const override {
return op->IsWaitingForScan();
}
void SetWaiting(TOperation::TPtr op) override {
op->SetWaitingForScanFlag();
}
void ResetWaiting(TOperation::TPtr op) override {
op->ResetWaitingForScanFlag();
}
bool Run(TOperation::TPtr op, TTransactionContext& txc, const TActorContext& ctx) override {
TActiveTransaction* tx = dynamic_cast<TActiveTransaction*>(op.Get());
Y_VERIFY_S(tx, "cannot cast operation of kind " << op->GetKind());
Y_VERIFY(tx->GetSchemeTx().HasBackup());
const auto& backup = tx->GetSchemeTx().GetBackup();
const ui64 tableId = backup.GetTableId();
Y_VERIFY(DataShard.GetUserTables().contains(tableId));
const ui32 localTableId = DataShard.GetUserTables().at(tableId)->LocalTid;
Y_VERIFY(txc.DB.GetScheme().GetTableInfo(localTableId));
auto* appData = AppData(ctx);
std::shared_ptr<::NKikimr::NDataShard::IExport> exp;
if (backup.HasYTSettings()) {
auto* exportFactory = appData->DataShardExportFactory;
if (exportFactory) {
const auto& settings = backup.GetYTSettings();
std::shared_ptr<IExport>(exportFactory->CreateExportToYt(settings.GetUseTypeV3())).swap(exp);
} else {
Abort(op, ctx, "Exports to YT are disabled");
return false;
}
} else if (backup.HasS3Settings()) {
auto* exportFactory = appData->DataShardExportFactory;
if (exportFactory) {
std::shared_ptr<IExport>(exportFactory->CreateExportToS3()).swap(exp);
} else {
Abort(op, ctx, "Exports to S3 are disabled");
return false;
}
} else {
Abort(op, ctx, "Unsupported backup task");
return false;
}
const auto& columns = DataShard.GetUserTables().at(tableId)->Columns;
const auto& scanSettings = backup.GetScanSettings();
const ui64 rowsLimit = scanSettings.GetRowsBatchSize() ? scanSettings.GetRowsBatchSize() : Max<ui64>();
const ui64 bytesLimit = scanSettings.GetBytesBatchSize();
auto createUploader = [self = DataShard.SelfId(), txId = op->GetTxId(), columns, backup, exp]() {
return exp->CreateUploader(self, txId, columns, backup);
};
THolder<IBuffer> buffer{exp->CreateBuffer(columns, rowsLimit, bytesLimit)};
THolder<NTable::IScan> scan{CreateExportScan(std::move(buffer), createUploader)};
const auto& taskName = appData->DataShardConfig.GetBackupTaskName();
const auto taskPrio = appData->DataShardConfig.GetBackupTaskPriority();
ui64 readAheadLo = appData->DataShardConfig.GetBackupReadAheadLo();
if (ui64 readAheadLoOverride = DataShard.GetBackupReadAheadLoOverride(); readAheadLoOverride > 0) {
readAheadLo = readAheadLoOverride;
}
ui64 readAheadHi = appData->DataShardConfig.GetBackupReadAheadHi();
if (ui64 readAheadHiOverride = DataShard.GetBackupReadAheadHiOverride(); readAheadHiOverride > 0) {
readAheadHi = readAheadHiOverride;
}
tx->SetScanTask(DataShard.QueueScan(localTableId, scan.Release(), op->GetTxId(),
TScanOptions()
.SetResourceBroker(taskName, taskPrio)
.SetReadAhead(readAheadLo, readAheadHi)
.SetReadPrio(TScanOptions::EReadPrio::Low)
));
return true;
}
bool HasResult(TOperation::TPtr op) const override {
return op->HasScanResult();
}
void ProcessResult(TOperation::TPtr op, const TActorContext&) override {
TActiveTransaction* tx = dynamic_cast<TActiveTransaction*>(op.Get());
Y_VERIFY_S(tx, "cannot cast operation of kind " << op->GetKind());
auto* result = CheckedCast<TExportScanProduct*>(op->ScanResult().Get());
auto* schemeOp = DataShard.FindSchemaTx(op->GetTxId());
schemeOp->Success = result->Success;
schemeOp->Error = std::move(result->Error);
schemeOp->BytesProcessed = result->BytesRead;
schemeOp->RowsProcessed = result->RowsRead;
op->SetScanResult(nullptr);
tx->SetScanTask(0);
}
void Cancel(TActiveTransaction* tx, const TActorContext&) override {
if (!tx->GetScanTask()) {
return;
}
const ui64 tableId = tx->GetSchemeTx().GetBackup().GetTableId();
Y_VERIFY(DataShard.GetUserTables().contains(tableId));
const ui32 localTableId = DataShard.GetUserTables().at(tableId)->LocalTid;
DataShard.CancelScan(localTableId, tx->GetScanTask());
tx->SetScanTask(0);
}
public:
TBackupUnit(TDataShard& self, TPipeline& pipeline)
: TBase(EExecutionUnitKind::Backup, self, pipeline)
{
}
}; // TBackupUnit
THolder<TExecutionUnit> CreateBackupUnit(TDataShard& self, TPipeline& pipeline) {
return THolder(new TBackupUnit(self, pipeline));
}
} // NDataShard
} // NKikimr
|