From 2c40b3c8dcd8da649e20bd2fc0e9e80532672f78 Mon Sep 17 00:00:00 2001 From: Innokentii Mokin Date: Tue, 11 Jun 2024 08:42:31 +0300 Subject: [backups] refactor change sender (#5354) --- ydb/core/change_exchange/change_exchange.cpp | 17 +- ydb/core/change_exchange/change_exchange.h | 31 +- ydb/core/change_exchange/change_record.h | 24 +- .../change_exchange/change_sender_common_ops.cpp | 641 ------------------- .../change_exchange/change_sender_common_ops.h | 698 +++++++++++++++++++-- ydb/core/change_exchange/change_sender_resolver.h | 24 + ydb/core/change_exchange/ya.make | 1 - ydb/core/scheme/scheme_tabledefs.h | 12 + ydb/core/tx/datashard/cdc_stream_heartbeat.cpp | 2 +- ydb/core/tx/datashard/cdc_stream_scan.cpp | 2 +- ydb/core/tx/datashard/change_collector.cpp | 2 +- ydb/core/tx/datashard/change_record.h | 75 +++ .../tx/datashard/change_sender_async_index.cpp | 46 +- ydb/core/tx/datashard/change_sender_cdc_stream.cpp | 126 ++-- ydb/core/tx/datashard/datashard_change_sending.cpp | 8 +- .../tx/datashard/datashard_ut_change_collector.cpp | 8 +- .../tx/replication/service/json_change_record.h | 58 ++ ydb/core/tx/replication/service/table_writer.cpp | 49 +- 18 files changed, 942 insertions(+), 882 deletions(-) delete mode 100644 ydb/core/change_exchange/change_sender_common_ops.cpp create mode 100644 ydb/core/change_exchange/change_sender_resolver.h diff --git a/ydb/core/change_exchange/change_exchange.cpp b/ydb/core/change_exchange/change_exchange.cpp index 92df9b3cea7..dfb46f1173b 100644 --- a/ydb/core/change_exchange/change_exchange.cpp +++ b/ydb/core/change_exchange/change_exchange.cpp @@ -88,21 +88,26 @@ TString TEvChangeExchange::TEvRemoveRecords::ToString() const { << " }"; } -/// TEvRecords -TEvChangeExchange::TEvRecords::TEvRecords(const TVector& records) +// TEvRecords +TEvChangeExchange::TEvRecords::TEvRecords(const TChangeRecordVector& records) : Records(records) { } -TEvChangeExchange::TEvRecords::TEvRecords(TVector&& records) +TEvChangeExchange::TEvRecords::TEvRecords(TChangeRecordVector&& records) : Records(std::move(records)) { } + TString TEvChangeExchange::TEvRecords::ToString() const { - return TStringBuilder() << ToStringHeader() << " {" - << " Records [" << JoinSeq(",", Records) << "]" - << " }"; + return std::visit( + [&](auto& records){ + return TStringBuilder() << ToStringHeader() << " {" + << " Records " << ((TBaseChangeRecordContainer*)records.get())->Out() + << " }"; + }, + Records); } /// TEvForgetRecords diff --git a/ydb/core/change_exchange/change_exchange.h b/ydb/core/change_exchange/change_exchange.h index 9cd000b4c23..2b42e9b90f9 100644 --- a/ydb/core/change_exchange/change_exchange.h +++ b/ydb/core/change_exchange/change_exchange.h @@ -8,8 +8,33 @@ #include +namespace NKikimr { + +namespace NDataShard { + class TChangeRecord; +} + +namespace NReplication::NService { + class TChangeRecord; +} + +struct TBaseChangeRecordContainer { + virtual ~TBaseChangeRecordContainer() = default; + virtual TString Out() = 0; +}; + +template +struct TChangeRecordContainer {}; + +} + namespace NKikimr::NChangeExchange { +using TChangeRecordVector = std::variant< + std::shared_ptr>, + std::shared_ptr> +>; + struct TEvChangeExchange { enum EEv { // Enqueue for sending @@ -73,10 +98,10 @@ struct TEvChangeExchange { }; struct TEvRecords: public TEventLocal { - TVector Records; + TChangeRecordVector Records; - explicit TEvRecords(const TVector& records); - explicit TEvRecords(TVector&& records); + explicit TEvRecords(const TChangeRecordVector& records); + explicit TEvRecords(TChangeRecordVector&& records); TString ToString() const override; }; diff --git a/ydb/core/change_exchange/change_record.h b/ydb/core/change_exchange/change_record.h index 88959dd0dc7..e47175afa20 100644 --- a/ydb/core/change_exchange/change_record.h +++ b/ydb/core/change_exchange/change_record.h @@ -6,6 +6,8 @@ namespace NKikimr::NChangeExchange { +class IChangeSenderResolver; + class IChangeRecord: public TThrRefBase { public: using TPtr = TIntrusivePtr; @@ -34,11 +36,7 @@ public: virtual TString ToString() const = 0; virtual void Out(IOutputStream& out) const = 0; - template - T* Get() { - return dynamic_cast(this); - } - + virtual ui64 ResolvePartitionId(IChangeSenderResolver* const resolver) const = 0; }; // IChangeRecord template class TChangeRecordBuilder; @@ -77,13 +75,11 @@ protected: public: TChangeRecordBuilder() : Record(MakeIntrusive()) - { - } + {} - explicit TChangeRecordBuilder(IChangeRecord::TPtr record) + explicit TChangeRecordBuilder(TIntrusivePtr record) : Record(std::move(record)) - { - } + {} TSelf& WithOrder(ui64 order) { GetRecord()->Order = order; @@ -105,12 +101,12 @@ public: return static_cast(*this); } - IChangeRecord::TPtr Build() { + TIntrusivePtr Build() { return Record; } protected: - IChangeRecord::TPtr Record; + TIntrusivePtr Record; }; // TChangeRecordBuilder @@ -119,3 +115,7 @@ protected: Y_DECLARE_OUT_SPEC(inline, NKikimr::NChangeExchange::IChangeRecord::TPtr, out, value) { return value->Out(out); } + +Y_DECLARE_OUT_SPEC(inline, NKikimr::NChangeExchange::IChangeRecord*, out, value) { + return value->Out(out); +} diff --git a/ydb/core/change_exchange/change_sender_common_ops.cpp b/ydb/core/change_exchange/change_sender_common_ops.cpp deleted file mode 100644 index 8b10355730a..00000000000 --- a/ydb/core/change_exchange/change_sender_common_ops.cpp +++ /dev/null @@ -1,641 +0,0 @@ -#include "change_sender_common_ops.h" -#include "change_sender_monitoring.h" - -#include - -#include -#include - -#include -#include - -namespace NKikimr::NChangeExchange { - -void TBaseChangeSender::LazyCreateSender(THashMap& senders, ui64 partitionId) { - auto res = senders.emplace(partitionId, TSender{}); - Y_ABORT_UNLESS(res.second); - - for (const auto& [order, broadcast] : Broadcasting) { - if (AddBroadcastPartition(order, partitionId)) { - // re-enqueue record to send it in the correct order - Enqueued.insert(ReEnqueue(broadcast.Record)); - } - } -} - -void TBaseChangeSender::RegisterSender(ui64 partitionId) { - Y_ABORT_UNLESS(Senders.contains(partitionId)); - auto& sender = Senders.at(partitionId); - - Y_ABORT_UNLESS(!sender.ActorId); - sender.ActorId = ActorOps->RegisterWithSameMailbox(CreateSender(partitionId)); -} - -void TBaseChangeSender::CreateMissingSenders(const TVector& partitionIds) { - THashMap senders; - - for (const auto& partitionId : partitionIds) { - auto it = Senders.find(partitionId); - if (it != Senders.end()) { - senders.emplace(partitionId, std::move(it->second)); - Senders.erase(it); - } else { - LazyCreateSender(senders, partitionId); - } - } - - for (const auto& [partitionId, sender] : Senders) { - ReEnqueueRecords(sender); - ProcessBroadcasting(&TBaseChangeSender::RemoveBroadcastPartition, - partitionId, sender.Broadcasting); - if (sender.ActorId) { - ActorOps->Send(sender.ActorId, new TEvents::TEvPoisonPill()); - } - } - - Senders = std::move(senders); -} - -void TBaseChangeSender::RecreateSenders(const TVector& partitionIds) { - for (const auto& partitionId : partitionIds) { - LazyCreateSender(Senders, partitionId); - } -} - -void TBaseChangeSender::CreateSenders(const TVector& partitionIds, bool partitioningChanged) { - if (partitioningChanged) { - CreateMissingSenders(partitionIds); - } else { - RecreateSenders(GonePartitions); - } - - GonePartitions.clear(); - - if (!Enqueued || !RequestRecords()) { - SendRecords(); - } -} - -void TBaseChangeSender::KillSenders() { - for (const auto& [_, sender] : std::exchange(Senders, {})) { - if (sender.ActorId) { - ActorOps->Send(sender.ActorId, new TEvents::TEvPoisonPill()); - } - } -} - -void TBaseChangeSender::EnqueueRecords(TVector&& records) { - for (auto& record : records) { - Y_VERIFY_S(PathId == record.PathId, "Unexpected record's path id" - << ": expected# " << PathId - << ", got# " << record.PathId); - Enqueued.emplace(record.Order, record.BodySize); - } - - RequestRecords(); -} - -bool TBaseChangeSender::RequestRecords() { - if (!Enqueued) { - return false; - } - - auto it = Enqueued.begin(); - TVector records; - - bool exceeded = false; - while (it != Enqueued.end()) { - if (MemUsage && (MemUsage + it->BodySize) > MemLimit) { - if (!it->ReEnqueued || exceeded) { - break; - } - - exceeded = true; - } - - MemUsage += it->BodySize; - - records.emplace_back(it->Order, it->BodySize); - PendingBody.emplace(it->Order, it->BodySize); - it = Enqueued.erase(it); - } - - if (!records) { - return false; - } - - ActorOps->Send(GetChangeServer(), new TEvChangeExchange::TEvRequestRecords(std::move(records))); - return true; -} - -void TBaseChangeSender::ProcessRecords(TVector&& records) { - for (auto& record : records) { - auto it = PendingBody.find(record->GetOrder()); - if (it == PendingBody.end()) { - continue; - } - - if (it->BodySize != record->GetBody().size()) { - MemUsage -= it->BodySize; - MemUsage += record->GetBody().size(); - } - - if (record->IsBroadcast()) { - // assume that broadcast records are too small to affect memory consumption - MemUsage -= record->GetBody().size(); - } - - PendingSent.emplace(record->GetOrder(), std::move(record)); - PendingBody.erase(it); - } - - SendRecords(); -} - -void TBaseChangeSender::SendRecords() { - if (!Resolver->IsResolved()) { - return; - } - - if (!PendingSent) { - return; - } - - auto it = PendingSent.begin(); - THashSet sendTo; - THashSet registrations; - bool needToResolve = false; - - while (it != PendingSent.end()) { - if (Enqueued && Enqueued.begin()->Order <= it->first) { - break; - } - - if (PendingBody && PendingBody.begin()->Order <= it->first) { - break; - } - - if (!it->second->IsBroadcast()) { - const ui64 partitionId = Resolver->GetPartitionId(it->second); - if (!Senders.contains(partitionId)) { - needToResolve = true; - ++it; - continue; - } - - auto& sender = Senders.at(partitionId); - sender.Prepared.push_back(std::move(it->second)); - if (!sender.ActorId) { - Y_ABORT_UNLESS(!sender.Ready); - registrations.insert(partitionId); - } - if (sender.Ready) { - sendTo.insert(partitionId); - } - } else { - auto& broadcast = EnsureBroadcast(it->second); - EraseNodesIf(broadcast.PendingPartitions, [&](ui64 partitionId) { - if (Senders.contains(partitionId)) { - auto& sender = Senders.at(partitionId); - sender.Prepared.push_back(std::move(it->second)); - if (!sender.ActorId) { - Y_ABORT_UNLESS(!sender.Ready); - registrations.insert(partitionId); - } - if (sender.Ready) { - sendTo.insert(partitionId); - } - - return true; - } - - return false; - }); - } - - it = PendingSent.erase(it); - } - - for (const auto partitionId : registrations) { - RegisterSender(partitionId); - } - - for (const auto partitionId : sendTo) { - SendPreparedRecords(partitionId); - } - - if (needToResolve && !Resolver->IsResolving()) { - Resolver->Resolve(); - } - - RequestRecords(); -} - -void TBaseChangeSender::ForgetRecords(TVector&& records) { - for (const auto& record : records) { - auto it = PendingBody.find(record); - if (it == PendingBody.end()) { - continue; - } - - MemUsage -= it->BodySize; - PendingBody.erase(it); - } - - RequestRecords(); -} - -void TBaseChangeSender::OnReady(ui64 partitionId) { - auto it = Senders.find(partitionId); - if (it == Senders.end()) { - return; - } - - auto& sender = it->second; - sender.Ready = true; - - if (sender.Pending) { - RemoveRecords(std::exchange(sender.Pending, {})); - } - - if (sender.Broadcasting) { - ProcessBroadcasting(&TBaseChangeSender::CompleteBroadcastPartition, - partitionId, std::exchange(sender.Broadcasting, {})); - } - - if (sender.Prepared) { - SendPreparedRecords(partitionId); - } - - RequestRecords(); -} - -void TBaseChangeSender::OnGone(ui64 partitionId) { - auto it = Senders.find(partitionId); - if (it == Senders.end()) { - return; - } - - ReEnqueueRecords(it->second); - Senders.erase(it); - GonePartitions.push_back(partitionId); - - if (Resolver->IsResolving()) { - return; - } - - Resolver->Resolve(); -} - -void TBaseChangeSender::SendPreparedRecords(ui64 partitionId) { - Y_ABORT_UNLESS(Senders.contains(partitionId)); - auto& sender = Senders.at(partitionId); - - Y_ABORT_UNLESS(sender.Ready); - sender.Ready = false; - - sender.Pending.reserve(sender.Prepared.size()); - for (const auto& record : sender.Prepared) { - if (!record->IsBroadcast()) { - sender.Pending.emplace_back(record->GetOrder(), record->GetBody().size()); - MemUsage -= record->GetBody().size(); - } else { - sender.Broadcasting.push_back(record->GetOrder()); - } - } - - Y_ABORT_UNLESS(sender.ActorId); - ActorOps->Send(sender.ActorId, new TEvChangeExchange::TEvRecords(std::exchange(sender.Prepared, {}))); -} - -void TBaseChangeSender::ReEnqueueRecords(const TSender& sender) { - for (const auto& record : sender.Pending) { - Enqueued.insert(ReEnqueue(record)); - } - - for (const auto& record : sender.Prepared) { - if (!record->IsBroadcast()) { - Enqueued.insert(ReEnqueue(record->GetOrder(), record->GetBody().size())); - MemUsage -= record->GetBody().size(); - } - } -} - -TBaseChangeSender::TBroadcast& TBaseChangeSender::EnsureBroadcast(IChangeRecord::TPtr record) { - Y_ABORT_UNLESS(record->IsBroadcast()); - - auto it = Broadcasting.find(record->GetOrder()); - if (it != Broadcasting.end()) { - return it->second; - } - - THashSet partitionIds; - for (const auto& [partitionId, _] : Senders) { - partitionIds.insert(partitionId); - } - for (const auto partitionId : GonePartitions) { - partitionIds.insert(partitionId); - } - - auto res = Broadcasting.emplace(record->GetOrder(), TBroadcast{ - .Record = {record->GetOrder(), record->GetBody().size()}, - .Partitions = partitionIds, - .PendingPartitions = partitionIds, - }); - - return res.first->second; -} - -bool TBaseChangeSender::AddBroadcastPartition(ui64 order, ui64 partitionId) { - auto it = Broadcasting.find(order); - Y_ABORT_UNLESS(it != Broadcasting.end()); - - auto& broadcast = it->second; - if (broadcast.Partitions.contains(partitionId)) { - return false; - } - - broadcast.Partitions.insert(partitionId); - broadcast.PendingPartitions.insert(partitionId); - - return true; -} - -bool TBaseChangeSender::RemoveBroadcastPartition(ui64 order, ui64 partitionId) { - auto it = Broadcasting.find(order); - Y_ABORT_UNLESS(it != Broadcasting.end()); - - auto& broadcast = it->second; - broadcast.Partitions.erase(partitionId); - broadcast.PendingPartitions.erase(partitionId); - broadcast.CompletedPartitions.erase(partitionId); - - return MaybeCompleteBroadcast(order); -} - -bool TBaseChangeSender::CompleteBroadcastPartition(ui64 order, ui64 partitionId) { - auto it = Broadcasting.find(order); - Y_ABORT_UNLESS(it != Broadcasting.end()); - - auto& broadcast = it->second; - broadcast.CompletedPartitions.insert(partitionId); - - return MaybeCompleteBroadcast(order); -} - -bool TBaseChangeSender::MaybeCompleteBroadcast(ui64 order) { - auto it = Broadcasting.find(order); - Y_ABORT_UNLESS(it != Broadcasting.end()); - - auto& broadcast = it->second; - if (broadcast.PendingPartitions || broadcast.Partitions.size() != broadcast.CompletedPartitions.size()) { - return false; - } - - Broadcasting.erase(it); - return true; -} - -void TBaseChangeSender::ProcessBroadcasting(std::function f, - ui64 partitionId, const TVector& broadcasting) -{ - TVector remove; - for (const auto order : broadcasting) { - if (std::invoke(f, this, order, partitionId)) { - remove.push_back(order); - } - } - - if (remove) { - RemoveRecords(std::move(remove)); - } -} - -void TBaseChangeSender::RemoveRecords() { - THashSet remove; - - for (const auto& record : std::exchange(Enqueued, {})) { - remove.insert(record.Order); - } - - for (const auto& record : std::exchange(PendingBody, {})) { - remove.insert(record.Order); - } - - for (const auto& [order, _] : std::exchange(PendingSent, {})) { - remove.insert(order); - } - - for (const auto& [order, _] : std::exchange(Broadcasting, {})) { - remove.insert(order); - } - - for (const auto& [_, sender] : Senders) { - for (const auto& record : sender.Pending) { - remove.insert(record.Order); - } - for (const auto& record : sender.Prepared) { - remove.insert(record->GetOrder()); - } - } - - if (remove) { - RemoveRecords(TVector(remove.begin(), remove.end())); - } -} - -TBaseChangeSender::TBaseChangeSender(IActorOps* actorOps, IChangeSenderResolver* resolver, const TPathId& pathId) - : ActorOps(actorOps) - , Resolver(resolver) - , PathId(pathId) - , MemLimit(192_KB) - , MemUsage(0) -{ -} - -void TBaseChangeSender::RenderHtmlPage(ui64 tabletId, NMon::TEvRemoteHttpInfo::TPtr& ev, - const TActorContext& ctx) -{ - const auto& cgi = ev->Get()->Cgi(); - if (const auto& str = cgi.Get("partitionId")) { - ui64 partitionId = 0; - if (TryFromString(str, partitionId)) { - auto it = Senders.find(partitionId); - if (it != Senders.end()) { - if (const auto& to = it->second.ActorId) { - ctx.Send(ev->Forward(to)); - } else { - ActorOps->Send(ev->Sender, new NMon::TEvRemoteHttpInfoRes(TStringBuilder() - << "Change sender '" << PathId << ":" << partitionId << "' is not running")); - } - } else { - ActorOps->Send(ev->Sender, new NMon::TEvRemoteBinaryInfoRes(NMonitoring::HTTPNOTFOUND)); - } - } else { - ActorOps->Send(ev->Sender, new NMon::TEvRemoteHttpInfoRes("Invalid partitionId")); - } - - return; - } - - TStringStream html; - - HTML(html) { - Header(html, "Change sender", tabletId); - - SimplePanel(html, "Info", [this](IOutputStream& html) { - HTML(html) { - DL_CLASS("dl-horizontal") { - TermDesc(html, "MemLimit", MemLimit); - TermDesc(html, "MemUsage", MemUsage); - } - } - }); - - SimplePanel(html, "Partition senders", [this, tabletId](IOutputStream& html) { - HTML(html) { - TABLE_CLASS("table table-hover") { - TABLEHEAD() { - TABLER() { - TABLEH() { html << "#"; } - TABLEH() { html << "PartitionId"; } - TABLEH() { html << "Ready"; } - TABLEH() { html << "Pending"; } - TABLEH() { html << "Prepared"; } - TABLEH() { html << "Broadcasting"; } - TABLEH() { html << "Actor"; } - } - } - TABLEBODY() { - ui32 i = 0; - for (const auto& [partitionId, sender] : Senders) { - TABLER() { - TABLED() { html << ++i; } - TABLED() { html << partitionId; } - TABLED() { html << sender.Ready; } - TABLED() { html << sender.Pending.size(); } - TABLED() { html << sender.Prepared.size(); } - TABLED() { html << sender.Broadcasting.size(); } - TABLED() { ActorLink(html, tabletId, PathId, partitionId); } - } - } - } - } - } - }); - - CollapsedPanel(html, "Enqueued", "enqueued", [this](IOutputStream& html) { - HTML(html) { - TABLE_CLASS("table table-hover") { - TABLEHEAD() { - TABLER() { - TABLEH() { html << "#"; } - TABLEH() { html << "Order"; } - TABLEH() { html << "BodySize"; } - } - } - TABLEBODY() { - ui32 i = 0; - for (const auto& record : Enqueued) { - TABLER() { - TABLED() { html << ++i; } - TABLED() { html << record.Order; } - TABLED() { html << record.BodySize; } - } - } - } - } - } - }); - - CollapsedPanel(html, "PendingBody", "pendingBody", [this](IOutputStream& html) { - HTML(html) { - TABLE_CLASS("table table-hover") { - TABLEHEAD() { - TABLER() { - TABLEH() { html << "#"; } - TABLEH() { html << "Order"; } - TABLEH() { html << "BodySize"; } - } - } - TABLEBODY() { - ui32 i = 0; - for (const auto& record : PendingBody) { - TABLER() { - TABLED() { html << ++i; } - TABLED() { html << record.Order; } - TABLED() { html << record.BodySize; } - } - } - } - } - } - }); - - CollapsedPanel(html, "PendingSent", "pendingSent", [this](IOutputStream& html) { - HTML(html) { - TABLE_CLASS("table table-hover") { - TABLEHEAD() { - TABLER() { - TABLEH() { html << "#"; } - TABLEH() { html << "Order"; } - TABLEH() { html << "Group"; } - TABLEH() { html << "Step"; } - TABLEH() { html << "TxId"; } - TABLEH() { html << "Kind"; } - TABLEH() { html << "Source"; } - } - } - TABLEBODY() { - ui32 i = 0; - for (const auto& [order, record] : PendingSent) { - TABLER() { - TABLED() { html << ++i; } - TABLED() { html << order; } - TABLED() { html << record->GetGroup(); } - TABLED() { html << record->GetStep(); } - TABLED() { html << record->GetTxId(); } - TABLED() { html << record->GetKind(); } - TABLED() { html << record->GetSource(); } - } - } - } - } - } - }); - - CollapsedPanel(html, "Broadcasting", "broadcasting", [this](IOutputStream& html) { - HTML(html) { - TABLE_CLASS("table table-hover") { - TABLEHEAD() { - TABLER() { - TABLEH() { html << "#"; } - TABLEH() { html << "Order"; } - TABLEH() { html << "BodySize"; } - TABLEH() { html << "Partitions"; } - TABLEH() { html << "PendingPartitions"; } - TABLEH() { html << "CompletedPartitions"; } - } - } - TABLEBODY() { - ui32 i = 0; - for (const auto& [order, broadcast] : Broadcasting) { - TABLER() { - TABLED() { html << ++i; } - TABLED() { html << order; } - TABLED() { html << broadcast.Record.BodySize; } - TABLED() { html << broadcast.Partitions.size(); } - TABLED() { html << broadcast.PendingPartitions.size(); } - TABLED() { html << broadcast.CompletedPartitions.size(); } - } - } - } - } - } - }); - } - - ActorOps->Send(ev->Sender, new NMon::TEvRemoteHttpInfoRes(html.Str())); -} - -} diff --git a/ydb/core/change_exchange/change_sender_common_ops.h b/ydb/core/change_exchange/change_sender_common_ops.h index 504afbde9da..0b684b30705 100644 --- a/ydb/core/change_exchange/change_sender_common_ops.h +++ b/ydb/core/change_exchange/change_sender_common_ops.h @@ -1,15 +1,24 @@ #pragma once #include "change_exchange.h" +#include "change_sender_resolver.h" + +#include #include #include +#include + +#include +#include #include #include #include #include +#include + namespace NKikimr::NChangeExchange { struct TEvChangeExchangePrivate { @@ -56,36 +65,17 @@ struct TEvChangeExchangePrivate { }; // TEvChangeExchangePrivate -class IChangeSender { -public: - virtual ~IChangeSender() = default; - - virtual TActorId GetChangeServer() const = 0; - - virtual void CreateSenders(const TVector& partitionIds, bool partitioningChanged = true) = 0; - virtual void KillSenders() = 0; - virtual IActor* CreateSender(ui64 partitionId) = 0; - virtual void RemoveRecords() = 0; - - virtual void EnqueueRecords(TVector&& records) = 0; - virtual void ProcessRecords(TVector&& records) = 0; - virtual void ForgetRecords(TVector&& records) = 0; - virtual void OnReady(ui64 partitionId) = 0; - virtual void OnGone(ui64 partitionId) = 0; -}; - -class IChangeSenderResolver { +class ISenderFactory { public: - virtual ~IChangeSenderResolver() = default; - - virtual void Resolve() = 0; - virtual bool IsResolving() const = 0; - virtual bool IsResolved() const = 0; - virtual ui64 GetPartitionId(IChangeRecord::TPtr record) const = 0; + virtual ~ISenderFactory() = default; + virtual IActor* CreateSender(ui64 partitionId) const = 0; }; -class TBaseChangeSender: public IChangeSender { +template +class TBaseChangeSender { using TIncompleteRecord = TEvChangeExchange::TEvRequestRecords::TRecordInfo; + // we need this to safely cast and call Out on a container + static_assert(std::derived_from, TBaseChangeRecordContainer>); struct TEnqueuedRecord: TIncompleteRecord { bool ReEnqueued = false; @@ -108,7 +98,7 @@ class TBaseChangeSender: public IChangeSender { TActorId ActorId; bool Ready = false; TVector Pending; - TVector Prepared; + TVector Prepared; TVector Broadcasting; }; @@ -119,24 +109,292 @@ class TBaseChangeSender: public IChangeSender { THashSet CompletedPartitions; }; - void LazyCreateSender(THashMap& senders, ui64 partitionId); - void RegisterSender(ui64 partitionId); - void CreateMissingSenders(const TVector& partitionIds); - void RecreateSenders(const TVector& partitionIds); + void LazyCreateSender(THashMap& senders, ui64 partitionId) { + auto res = senders.emplace(partitionId, TSender{}); + Y_ABORT_UNLESS(res.second); + + for (const auto& [order, broadcast] : Broadcasting) { + if (AddBroadcastPartition(order, partitionId)) { + // re-enqueue record to send it in the correct order + Enqueued.insert(ReEnqueue(broadcast.Record)); + } + } + } + + void RegisterSender(ui64 partitionId) { + Y_ABORT_UNLESS(Senders.contains(partitionId)); + auto& sender = Senders.at(partitionId); + + Y_ABORT_UNLESS(!sender.ActorId); + sender.ActorId = ActorOps->RegisterWithSameMailbox(SenderFactory->CreateSender(partitionId)); + } + + void CreateMissingSenders(const TVector& partitionIds) { + THashMap senders; + + for (const auto& partitionId : partitionIds) { + auto it = Senders.find(partitionId); + if (it != Senders.end()) { + senders.emplace(partitionId, std::move(it->second)); + Senders.erase(it); + } else { + LazyCreateSender(senders, partitionId); + } + } + + for (const auto& [partitionId, sender] : Senders) { + ReEnqueueRecords(sender); + ProcessBroadcasting(&TBaseChangeSender::RemoveBroadcastPartition, + partitionId, sender.Broadcasting); + if (sender.ActorId) { + ActorOps->Send(sender.ActorId, new TEvents::TEvPoisonPill()); + } + } + + Senders = std::move(senders); + } + + void RecreateSenders(const TVector& partitionIds) { + for (const auto& partitionId : partitionIds) { + LazyCreateSender(Senders, partitionId); + } + } + + bool RequestRecords() { + if (!Enqueued) { + return false; + } + + auto it = Enqueued.begin(); + TVector records; + + bool exceeded = false; + while (it != Enqueued.end()) { + if (MemUsage && (MemUsage + it->BodySize) > MemLimit) { + if (!it->ReEnqueued || exceeded) { + break; + } + + exceeded = true; + } + + MemUsage += it->BodySize; + + records.emplace_back(it->Order, it->BodySize); + PendingBody.emplace(it->Order, it->BodySize); + it = Enqueued.erase(it); + } + + if (!records) { + return false; + } + + ActorOps->Send(GetChangeServer(), new TEvChangeExchange::TEvRequestRecords(std::move(records))); + return true; + } + + void SendRecords() { + if (!Resolver->IsResolved()) { + return; + } + + if (!PendingSent) { + return; + } + + auto it = PendingSent.begin(); + THashSet sendTo; + THashSet registrations; + bool needToResolve = false; + + while (it != PendingSent.end()) { + if (Enqueued && Enqueued.begin()->Order <= it->first) { + break; + } + + if (PendingBody && PendingBody.begin()->Order <= it->first) { + break; + } + + if (!it->second->IsBroadcast()) { + const ui64 partitionId = it->second->ResolvePartitionId(Resolver); + if (!Senders.contains(partitionId)) { + needToResolve = true; + ++it; + continue; + } + + auto& sender = Senders.at(partitionId); + sender.Prepared.push_back(std::move(it->second)); + if (!sender.ActorId) { + Y_ABORT_UNLESS(!sender.Ready); + registrations.insert(partitionId); + } + if (sender.Ready) { + sendTo.insert(partitionId); + } + } else { + auto& broadcast = EnsureBroadcast(it->second); + EraseNodesIf(broadcast.PendingPartitions, [&](ui64 partitionId) { + if (Senders.contains(partitionId)) { + auto& sender = Senders.at(partitionId); + sender.Prepared.push_back(std::move(it->second)); + if (!sender.ActorId) { + Y_ABORT_UNLESS(!sender.Ready); + registrations.insert(partitionId); + } + if (sender.Ready) { + sendTo.insert(partitionId); + } + + return true; + } + + return false; + }); + } + + it = PendingSent.erase(it); + } + + for (const auto partitionId : registrations) { + RegisterSender(partitionId); + } + + for (const auto partitionId : sendTo) { + SendPreparedRecords(partitionId); + } + + if (needToResolve && !Resolver->IsResolving()) { + Resolver->Resolve(); + } + + RequestRecords(); + } + + void SendPreparedRecords(ui64 partitionId) { + Y_ABORT_UNLESS(Senders.contains(partitionId)); + auto& sender = Senders.at(partitionId); + + Y_ABORT_UNLESS(sender.Ready); + sender.Ready = false; + + sender.Pending.reserve(sender.Prepared.size()); + for (const auto& record : sender.Prepared) { + if (!record->IsBroadcast()) { + sender.Pending.emplace_back(record->GetOrder(), record->GetBody().size()); + MemUsage -= record->GetBody().size(); + } else { + sender.Broadcasting.push_back(record->GetOrder()); + } + } + + Y_ABORT_UNLESS(sender.ActorId); + ActorOps->Send(sender.ActorId, new TEvChangeExchange::TEvRecords(std::make_shared>(std::exchange(sender.Prepared, {})))); + } + + void ReEnqueueRecords(const TSender& sender) { + for (const auto& record : sender.Pending) { + Enqueued.insert(ReEnqueue(record)); + } + + for (const auto& record : sender.Prepared) { + if (!record->IsBroadcast()) { + Enqueued.insert(ReEnqueue(record->GetOrder(), record->GetBody().size())); + MemUsage -= record->GetBody().size(); + } + } + } + + TBroadcast& EnsureBroadcast(IChangeRecord::TPtr record) { + Y_ABORT_UNLESS(record->IsBroadcast()); + + auto it = Broadcasting.find(record->GetOrder()); + if (it != Broadcasting.end()) { + return it->second; + } + + THashSet partitionIds; + for (const auto& [partitionId, _] : Senders) { + partitionIds.insert(partitionId); + } + for (const auto partitionId : GonePartitions) { + partitionIds.insert(partitionId); + } + + auto res = Broadcasting.emplace(record->GetOrder(), TBroadcast{ + .Record = {record->GetOrder(), record->GetBody().size()}, + .Partitions = partitionIds, + .PendingPartitions = partitionIds, + }); + + return res.first->second; + } + + bool AddBroadcastPartition(ui64 order, ui64 partitionId) { + auto it = Broadcasting.find(order); + Y_ABORT_UNLESS(it != Broadcasting.end()); + + auto& broadcast = it->second; + if (broadcast.Partitions.contains(partitionId)) { + return false; + } + + broadcast.Partitions.insert(partitionId); + broadcast.PendingPartitions.insert(partitionId); + + return true; + } + + bool RemoveBroadcastPartition(ui64 order, ui64 partitionId) { + auto it = Broadcasting.find(order); + Y_ABORT_UNLESS(it != Broadcasting.end()); + + auto& broadcast = it->second; + broadcast.Partitions.erase(partitionId); + broadcast.PendingPartitions.erase(partitionId); + broadcast.CompletedPartitions.erase(partitionId); + + return MaybeCompleteBroadcast(order); + } + + bool CompleteBroadcastPartition(ui64 order, ui64 partitionId) { + auto it = Broadcasting.find(order); + Y_ABORT_UNLESS(it != Broadcasting.end()); - bool RequestRecords(); - void SendRecords(); + auto& broadcast = it->second; + broadcast.CompletedPartitions.insert(partitionId); - void SendPreparedRecords(ui64 partitionId); - void ReEnqueueRecords(const TSender& sender); + return MaybeCompleteBroadcast(order); + } + + bool MaybeCompleteBroadcast(ui64 order) { + auto it = Broadcasting.find(order); + Y_ABORT_UNLESS(it != Broadcasting.end()); + + auto& broadcast = it->second; + if (broadcast.PendingPartitions || broadcast.Partitions.size() != broadcast.CompletedPartitions.size()) { + return false; + } + + Broadcasting.erase(it); + return true; + } - TBroadcast& EnsureBroadcast(IChangeRecord::TPtr record); - bool AddBroadcastPartition(ui64 order, ui64 partitionId); - bool RemoveBroadcastPartition(ui64 order, ui64 partitionId); - bool CompleteBroadcastPartition(ui64 order, ui64 partitionId); - bool MaybeCompleteBroadcast(ui64 order); void ProcessBroadcasting(std::function f, - ui64 partitionId, const TVector& broadcasting); + ui64 partitionId, const TVector& broadcasting) + { + TVector remove; + for (const auto order : broadcasting) { + if (std::invoke(f, this, order, partitionId)) { + remove.push_back(order); + } + } + + if (remove) { + RemoveRecords(std::move(remove)); + } + } protected: template @@ -154,25 +412,358 @@ protected: ActorOps->Send(GetChangeServer(), new TEvChangeExchange::TEvRemoveRecords(std::move(records))); } - void CreateSenders(const TVector& partitionIds, bool partitioningChanged = true) override; - void KillSenders() override; - void RemoveRecords() override; + TActorId GetChangeServer() const { return ChangeServer; } + void CreateSenders(const TVector& partitionIds, bool partitioningChanged = true) { + if (partitioningChanged) { + CreateMissingSenders(partitionIds); + } else { + RecreateSenders(GonePartitions); + } + + GonePartitions.clear(); + + if (!Enqueued || !RequestRecords()) { + SendRecords(); + } + } + + void KillSenders() { + for (const auto& [_, sender] : std::exchange(Senders, {})) { + if (sender.ActorId) { + ActorOps->Send(sender.ActorId, new TEvents::TEvPoisonPill()); + } + } + } + + void RemoveRecords() { + THashSet remove; - void EnqueueRecords(TVector&& records) override; - void ProcessRecords(TVector&& records) override; - void ForgetRecords(TVector&& records) override; - void OnReady(ui64 partitionId) override; - void OnGone(ui64 partitionId) override; + for (const auto& record : std::exchange(Enqueued, {})) { + remove.insert(record.Order); + } - explicit TBaseChangeSender(IActorOps* actorOps, IChangeSenderResolver* resolver, const TPathId& pathId); + for (const auto& record : std::exchange(PendingBody, {})) { + remove.insert(record.Order); + } - void RenderHtmlPage(ui64 tabletId, NMon::TEvRemoteHttpInfo::TPtr& ev, const TActorContext& ctx); + for (const auto& [order, _] : std::exchange(PendingSent, {})) { + remove.insert(order); + } + + for (const auto& [order, _] : std::exchange(Broadcasting, {})) { + remove.insert(order); + } + + for (const auto& [_, sender] : Senders) { + for (const auto& record : sender.Pending) { + remove.insert(record.Order); + } + + for (const auto& record : sender.Prepared) { + remove.insert(record->GetOrder()); + } + } + + if (remove) { + RemoveRecords(TVector(remove.begin(), remove.end())); + } + } + + void EnqueueRecords(TVector&& records) { + for (auto& record : records) { + Y_VERIFY_S(PathId == record.PathId, "Unexpected record's path id" + << ": expected# " << PathId + << ", got# " << record.PathId); + Enqueued.emplace(record.Order, record.BodySize); + } + + RequestRecords(); + } + + void ProcessRecords(TVector&& records) { + for (auto& record : records) { + auto it = PendingBody.find(record->GetOrder()); + if (it == PendingBody.end()) { + continue; + } + + if (it->BodySize != record->GetBody().size()) { + MemUsage -= it->BodySize; + MemUsage += record->GetBody().size(); + } + + if (record->IsBroadcast()) { + // assume that broadcast records are too small to affect memory consumption + MemUsage -= record->GetBody().size(); + } + + PendingSent.emplace(record->GetOrder(), std::move(record)); + PendingBody.erase(it); + } + + SendRecords(); + } + + void ForgetRecords(TVector&& records) { + for (const auto& record : records) { + auto it = PendingBody.find(record); + if (it == PendingBody.end()) { + continue; + } + + MemUsage -= it->BodySize; + PendingBody.erase(it); + } + + RequestRecords(); + } + + void OnReady(ui64 partitionId) { + auto it = Senders.find(partitionId); + if (it == Senders.end()) { + return; + } + + auto& sender = it->second; + sender.Ready = true; + + if (sender.Pending) { + RemoveRecords(std::exchange(sender.Pending, {})); + } + + if (sender.Broadcasting) { + ProcessBroadcasting(&TBaseChangeSender::CompleteBroadcastPartition, + partitionId, std::exchange(sender.Broadcasting, {})); + } + + if (sender.Prepared) { + SendPreparedRecords(partitionId); + } + + RequestRecords(); + } + + void OnGone(ui64 partitionId) { + auto it = Senders.find(partitionId); + if (it == Senders.end()) { + return; + } + + ReEnqueueRecords(it->second); + Senders.erase(it); + GonePartitions.push_back(partitionId); + + if (Resolver->IsResolving()) { + return; + } + + Resolver->Resolve(); + } + + explicit TBaseChangeSender( + IActorOps* const actorOps, + IChangeSenderResolver* const resolver, + ISenderFactory* const senderFactory, + const TActorId changeServer, + const TPathId& pathId) + : ActorOps(actorOps) + , Resolver(resolver) + , SenderFactory(senderFactory) + , ChangeServer(changeServer) + , PathId(pathId) + , MemLimit(192_KB) + , MemUsage(0) + {} + + void RenderHtmlPage(ui64 tabletId, NMon::TEvRemoteHttpInfo::TPtr& ev, const TActorContext& ctx) { + const auto& cgi = ev->Get()->Cgi(); + if (const auto& str = cgi.Get("partitionId")) { + ui64 partitionId = 0; + if (TryFromString(str, partitionId)) { + auto it = Senders.find(partitionId); + if (it != Senders.end()) { + if (const auto& to = it->second.ActorId) { + ctx.Send(ev->Forward(to)); + } else { + ActorOps->Send(ev->Sender, new NMon::TEvRemoteHttpInfoRes(TStringBuilder() + << "Change sender '" << PathId << ":" << partitionId << "' is not running")); + } + } else { + ActorOps->Send(ev->Sender, new NMon::TEvRemoteBinaryInfoRes(NMonitoring::HTTPNOTFOUND)); + } + } else { + ActorOps->Send(ev->Sender, new NMon::TEvRemoteHttpInfoRes("Invalid partitionId")); + } + + return; + } + + TStringStream html; + + HTML(html) { + Header(html, "Change sender", tabletId); + + SimplePanel(html, "Info", [this](IOutputStream& html) { + HTML(html) { + DL_CLASS("dl-horizontal") { + TermDesc(html, "MemLimit", MemLimit); + TermDesc(html, "MemUsage", MemUsage); + } + } + }); + + SimplePanel(html, "Partition senders", [this, tabletId](IOutputStream& html) { + HTML(html) { + TABLE_CLASS("table table-hover") { + TABLEHEAD() { + TABLER() { + TABLEH() { html << "#"; } + TABLEH() { html << "PartitionId"; } + TABLEH() { html << "Ready"; } + TABLEH() { html << "Pending"; } + TABLEH() { html << "Prepared"; } + TABLEH() { html << "Broadcasting"; } + TABLEH() { html << "Actor"; } + } + } + TABLEBODY() { + ui32 i = 0; + for (const auto& [partitionId, sender] : Senders) { + TABLER() { + TABLED() { html << ++i; } + TABLED() { html << partitionId; } + TABLED() { html << sender.Ready; } + TABLED() { html << sender.Pending.size(); } + TABLED() { html << sender.Prepared.size(); } + TABLED() { html << sender.Broadcasting.size(); } + TABLED() { ActorLink(html, tabletId, PathId, partitionId); } + } + } + } + } + } + }); + + CollapsedPanel(html, "Enqueued", "enqueued", [this](IOutputStream& html) { + HTML(html) { + TABLE_CLASS("table table-hover") { + TABLEHEAD() { + TABLER() { + TABLEH() { html << "#"; } + TABLEH() { html << "Order"; } + TABLEH() { html << "BodySize"; } + } + } + TABLEBODY() { + ui32 i = 0; + for (const auto& record : Enqueued) { + TABLER() { + TABLED() { html << ++i; } + TABLED() { html << record.Order; } + TABLED() { html << record.BodySize; } + } + } + } + } + } + }); + + CollapsedPanel(html, "PendingBody", "pendingBody", [this](IOutputStream& html) { + HTML(html) { + TABLE_CLASS("table table-hover") { + TABLEHEAD() { + TABLER() { + TABLEH() { html << "#"; } + TABLEH() { html << "Order"; } + TABLEH() { html << "BodySize"; } + } + } + TABLEBODY() { + ui32 i = 0; + for (const auto& record : PendingBody) { + TABLER() { + TABLED() { html << ++i; } + TABLED() { html << record.Order; } + TABLED() { html << record.BodySize; } + } + } + } + } + } + }); + + CollapsedPanel(html, "PendingSent", "pendingSent", [this](IOutputStream& html) { + HTML(html) { + TABLE_CLASS("table table-hover") { + TABLEHEAD() { + TABLER() { + TABLEH() { html << "#"; } + TABLEH() { html << "Order"; } + TABLEH() { html << "Group"; } + TABLEH() { html << "Step"; } + TABLEH() { html << "TxId"; } + TABLEH() { html << "Kind"; } + TABLEH() { html << "Source"; } + } + } + TABLEBODY() { + ui32 i = 0; + for (const auto& [order, record] : PendingSent) { + TABLER() { + TABLED() { html << ++i; } + TABLED() { html << order; } + TABLED() { html << record->GetGroup(); } + TABLED() { html << record->GetStep(); } + TABLED() { html << record->GetTxId(); } + TABLED() { html << record->GetKind(); } + TABLED() { html << record->GetSource(); } + } + } + } + } + } + }); + + CollapsedPanel(html, "Broadcasting", "broadcasting", [this](IOutputStream& html) { + HTML(html) { + TABLE_CLASS("table table-hover") { + TABLEHEAD() { + TABLER() { + TABLEH() { html << "#"; } + TABLEH() { html << "Order"; } + TABLEH() { html << "BodySize"; } + TABLEH() { html << "Partitions"; } + TABLEH() { html << "PendingPartitions"; } + TABLEH() { html << "CompletedPartitions"; } + } + } + TABLEBODY() { + ui32 i = 0; + for (const auto& [order, broadcast] : Broadcasting) { + TABLER() { + TABLED() { html << ++i; } + TABLED() { html << order; } + TABLED() { html << broadcast.Record.BodySize; } + TABLED() { html << broadcast.Partitions.size(); } + TABLED() { html << broadcast.PendingPartitions.size(); } + TABLED() { html << broadcast.CompletedPartitions.size(); } + } + } + } + } + } + }); + } + + ActorOps->Send(ev->Sender, new NMon::TEvRemoteHttpInfoRes(html.Str())); + } private: IActorOps* const ActorOps; IChangeSenderResolver* const Resolver; - + ISenderFactory* const SenderFactory; protected: + TActorId ChangeServer; const TPathId PathId; private: @@ -182,11 +773,10 @@ private: THashMap Senders; // ui64 is partition id TSet Enqueued; TSet PendingBody; - TMap PendingSent; // ui64 is order + TMap PendingSent; // ui64 is order THashMap Broadcasting; // ui64 is order TVector GonePartitions; - }; // TBaseChangeSender } diff --git a/ydb/core/change_exchange/change_sender_resolver.h b/ydb/core/change_exchange/change_sender_resolver.h new file mode 100644 index 00000000000..71cb67605a9 --- /dev/null +++ b/ydb/core/change_exchange/change_sender_resolver.h @@ -0,0 +1,24 @@ +#pragma once + +#include +#include +#include + +#include + +namespace NKikimr::NChangeExchange { + +class IChangeSenderResolver { +public: + virtual ~IChangeSenderResolver() = default; + + virtual void Resolve() = 0; + virtual bool IsResolving() const = 0; + virtual bool IsResolved() const = 0; + + virtual const TVector& GetPartitions() const = 0; + virtual const TVector& GetSchema() const = 0; + virtual NKikimrSchemeOp::ECdcStreamFormat GetStreamFormat() const = 0; +}; + +} // NKikimr::NChangeExchange diff --git a/ydb/core/change_exchange/ya.make b/ydb/core/change_exchange/ya.make index 6e82179d832..b95ab217844 100644 --- a/ydb/core/change_exchange/ya.make +++ b/ydb/core/change_exchange/ya.make @@ -3,7 +3,6 @@ LIBRARY() SRCS( change_exchange.cpp change_record.cpp - change_sender_common_ops.cpp change_sender_monitoring.cpp ) diff --git a/ydb/core/scheme/scheme_tabledefs.h b/ydb/core/scheme/scheme_tabledefs.h index bbe8f1cd99e..2e8dc6941b4 100644 --- a/ydb/core/scheme/scheme_tabledefs.h +++ b/ydb/core/scheme/scheme_tabledefs.h @@ -712,6 +712,18 @@ public: , Status(EStatus::Unknown) , Partitioning(std::make_shared>()) {} + + static THolder CreateMiniKeyDesc(const TVector &keyColumnTypes) { + return THolder(new TKeyDesc(keyColumnTypes)); + } +private: + TKeyDesc(const TVector &keyColumnTypes) + : RowOperation(ERowOperation::Unknown) + , KeyColumnTypes(keyColumnTypes.begin(), keyColumnTypes.end()) + , Reverse(false) + , Status(EStatus::Unknown) + , Partitioning(std::make_shared>()) + {} }; struct TSystemColumnInfo { diff --git a/ydb/core/tx/datashard/cdc_stream_heartbeat.cpp b/ydb/core/tx/datashard/cdc_stream_heartbeat.cpp index 0560136f72a..6473bc62ba6 100644 --- a/ydb/core/tx/datashard/cdc_stream_heartbeat.cpp +++ b/ydb/core/tx/datashard/cdc_stream_heartbeat.cpp @@ -51,7 +51,7 @@ public: .WithSchemaVersion(0) // not used .Build(); - const auto& record = *recordPtr->Get(); + const auto& record = *recordPtr; Self->PersistChangeRecord(db, record); ChangeRecords.push_back(IDataShardChangeCollector::TChange{ diff --git a/ydb/core/tx/datashard/cdc_stream_scan.cpp b/ydb/core/tx/datashard/cdc_stream_scan.cpp index 4bae27518d4..944f9de86c2 100644 --- a/ydb/core/tx/datashard/cdc_stream_scan.cpp +++ b/ydb/core/tx/datashard/cdc_stream_scan.cpp @@ -314,7 +314,7 @@ public: .WithSource(TChangeRecord::ESource::InitialScan) .Build(); - const auto& record = *recordPtr->Get(); + const auto& record = *recordPtr; Self->PersistChangeRecord(db, record); ChangeRecords.push_back(IDataShardChangeCollector::TChange{ diff --git a/ydb/core/tx/datashard/change_collector.cpp b/ydb/core/tx/datashard/change_collector.cpp index c2f30163d44..56f4e673ba5 100644 --- a/ydb/core/tx/datashard/change_collector.cpp +++ b/ydb/core/tx/datashard/change_collector.cpp @@ -141,7 +141,7 @@ public: .WithBody(body.SerializeAsString()) .Build(); - const auto& record = *recordPtr->Get(); + const auto& record = *recordPtr; Self->PersistChangeRecord(db, record); if (record.GetLockId() == 0) { diff --git a/ydb/core/tx/datashard/change_record.h b/ydb/core/tx/datashard/change_record.h index 236f7d4ce46..0366fd2c77d 100644 --- a/ydb/core/tx/datashard/change_record.h +++ b/ydb/core/tx/datashard/change_record.h @@ -3,10 +3,14 @@ #include "datashard_user_table.h" #include +#include +#include #include #include +#include #include +#include namespace NKikimrChangeExchange { class TChangeRecord; @@ -20,6 +24,8 @@ class TChangeRecord: public NChangeExchange::TChangeRecordBase { friend class TChangeRecordBuilder; public: + using TPtr = TIntrusivePtr; + ui64 GetGroup() const override { return Group; } ui64 GetStep() const override { return Step; } ui64 GetTxId() const override { return TxId; } @@ -42,6 +48,50 @@ public: void Out(IOutputStream& out) const override; + ui64 ResolvePartitionId(NChangeExchange::IChangeSenderResolver* const resolver) const override { + const auto& partitions = resolver->GetPartitions(); + Y_ABORT_UNLESS(partitions); + const auto& schema = resolver->GetSchema(); + const auto streamFormat = resolver->GetStreamFormat(); + + switch (streamFormat) { + case NKikimrSchemeOp::ECdcStreamFormatProto: { + const auto range = TTableRange(GetKey()); + Y_ABORT_UNLESS(range.Point); + + const auto it = LowerBound( + partitions.cbegin(), partitions.cend(), true, + [&](const auto& partition, bool) { + Y_ABORT_UNLESS(partition.Range); + const int compares = CompareBorders( + partition.Range->EndKeyPrefix.GetCells(), range.From, + partition.Range->IsInclusive || partition.Range->IsPoint, + range.InclusiveFrom || range.Point, schema + ); + + return (compares < 0); + } + ); + + Y_ABORT_UNLESS(it != partitions.cend()); + return it->ShardId; + } + + case NKikimrSchemeOp::ECdcStreamFormatJson: + case NKikimrSchemeOp::ECdcStreamFormatDynamoDBStreamsJson: + case NKikimrSchemeOp::ECdcStreamFormatDebeziumJson: { + using namespace NKikimr::NDataStreams::V1; + const auto hashKey = HexBytesToDecimal(GetPartitionKey() /* MD5 */); + return ShardFromDecimal(hashKey, partitions.size()); + } + + default: { + Y_FAIL_S("Unknown format" + << ": format# " << static_cast(streamFormat)); + } + } + } + private: ui64 Group = 0; ui64 Step = 0; @@ -119,6 +169,31 @@ public: } +namespace NKikimr { + +template <> +struct TChangeRecordContainer + : public TBaseChangeRecordContainer +{ + TChangeRecordContainer() = default; + + explicit TChangeRecordContainer(TVector&& records) + : Records(std::move(records)) + {} + + TVector Records; + + TString Out() override { + return TStringBuilder() << "[" << JoinSeq(",", Records) << "]"; + } +}; + +} + Y_DECLARE_OUT_SPEC(inline, NKikimr::NDataShard::TChangeRecord, out, value) { return value.Out(out); } + +Y_DECLARE_OUT_SPEC(inline, NKikimr::NDataShard::TChangeRecord::TPtr, out, value) { + return value->Out(out); +} diff --git a/ydb/core/tx/datashard/change_sender_async_index.cpp b/ydb/core/tx/datashard/change_sender_async_index.cpp index c0630662a4e..38492b20728 100644 --- a/ydb/core/tx/datashard/change_sender_async_index.cpp +++ b/ydb/core/tx/datashard/change_sender_async_index.cpp @@ -128,8 +128,10 @@ class TAsyncIndexChangeSenderShard: public TActorBootstrappedRecord.SetOrigin(DataShard.TabletId); records->Record.SetGeneration(DataShard.Generation); - for (auto recordPtr : ev->Get()->Records) { - const auto& record = *recordPtr->Get(); + auto& evRecords = std::get>>(ev->Get()->Records)->Records; + + for (auto& recordPtr : evRecords) { + const auto& record = *recordPtr; if (record.GetOrder() <= LastRecordOrder) { continue; @@ -330,8 +332,9 @@ private: class TAsyncIndexChangeSenderMain : public TActorBootstrapped - , public NChangeExchange::TBaseChangeSender + , public NChangeExchange::TBaseChangeSender , public NChangeExchange::IChangeSenderResolver + , public NChangeExchange::ISenderFactory , private NSchemeCache::TSchemeCacheHelpers { TStringBuf GetLogPrefix() const { @@ -704,10 +707,6 @@ class TAsyncIndexChangeSenderMain return StateBase(ev); } - TActorId GetChangeServer() const override { - return DataShard.ActorId; - } - void Resolve() override { ResolveIndex(); } @@ -716,31 +715,11 @@ class TAsyncIndexChangeSenderMain return KeyDesc && KeyDesc->GetPartitions(); } - ui64 GetPartitionId(NChangeExchange::IChangeRecord::TPtr record) const override { - Y_ABORT_UNLESS(KeyDesc); - Y_ABORT_UNLESS(KeyDesc->GetPartitions()); - - const auto range = TTableRange(record->Get()->GetKey()); - Y_ABORT_UNLESS(range.Point); - - TVector::const_iterator it = LowerBound( - KeyDesc->GetPartitions().begin(), KeyDesc->GetPartitions().end(), true, - [&](const TKeyDesc::TPartitionInfo& partition, bool) { - const int compares = CompareBorders( - partition.Range->EndKeyPrefix.GetCells(), range.From, - partition.Range->IsInclusive || partition.Range->IsPoint, - range.InclusiveFrom || range.Point, KeyDesc->KeyColumnTypes - ); - - return (compares < 0); - } - ); - - Y_ABORT_UNLESS(it != KeyDesc->GetPartitions().end()); - return it->ShardId; // partition = shard - } + const TVector& GetPartitions() const override { return KeyDesc->GetPartitions(); } + const TVector& GetSchema() const override { return KeyDesc->KeyColumnTypes; } + NKikimrSchemeOp::ECdcStreamFormat GetStreamFormat() const override { return NKikimrSchemeOp::ECdcStreamFormatProto; } - IActor* CreateSender(ui64 partitionId) override { + IActor* CreateSender(ui64 partitionId) const override { return new TAsyncIndexChangeSenderShard(SelfId(), DataShard, partitionId, IndexTablePathId, TagMap); } @@ -751,7 +730,8 @@ class TAsyncIndexChangeSenderMain void Handle(NChangeExchange::TEvChangeExchange::TEvRecords::TPtr& ev) { LOG_D("Handle " << ev->Get()->ToString()); - ProcessRecords(std::move(ev->Get()->Records)); + auto& records = std::get>>(ev->Get()->Records)->Records; + ProcessRecords(std::move(records)); } void Handle(NChangeExchange::TEvChangeExchange::TEvForgetRecords::TPtr& ev) { @@ -798,7 +778,7 @@ public: explicit TAsyncIndexChangeSenderMain(const TDataShardId& dataShard, const TTableId& userTableId, const TPathId& indexPathId) : TActorBootstrapped() - , TBaseChangeSender(this, this, indexPathId) + , TBaseChangeSender(this, this, this, dataShard.ActorId, indexPathId) , DataShard(dataShard) , UserTableId(userTableId) , IndexTableVersion(0) diff --git a/ydb/core/tx/datashard/change_sender_cdc_stream.cpp b/ydb/core/tx/datashard/change_sender_cdc_stream.cpp index 55cf00095e5..5300357c24c 100644 --- a/ydb/core/tx/datashard/change_sender_cdc_stream.cpp +++ b/ydb/core/tx/datashard/change_sender_cdc_stream.cpp @@ -95,8 +95,9 @@ class TCdcChangeSenderPartition: public TActorBootstrappedGet()->ToString()); NKikimrClient::TPersQueueRequest request; - for (auto recordPtr : ev->Get()->Records) { - const auto& record = *recordPtr->Get(); + auto& records = std::get>>(ev->Get()->Records)->Records; + for (auto recordPtr : records) { + const auto& record = *recordPtr; if (record.GetSeqNo() <= MaxSeqNo) { continue; @@ -294,8 +295,9 @@ private: class TCdcChangeSenderMain : public TActorBootstrapped - , public NChangeExchange::TBaseChangeSender + , public NChangeExchange::TBaseChangeSender , public NChangeExchange::IChangeSenderResolver + , public NChangeExchange::ISenderFactory , private NSchemeCache::TSchemeCacheHelpers { struct TPQPartitionInfo { @@ -337,31 +339,6 @@ class TCdcChangeSenderMain }; // TPQPartitionInfo - struct TKeyDesc { - struct TPartitionInfo { - ui32 PartitionId; - ui64 ShardId; - TSerializedCellVec EndKeyPrefix; - // just a hint - static constexpr bool IsInclusive = false; - static constexpr bool IsPoint = false; - - explicit TPartitionInfo(const TPQPartitionInfo& info) - : PartitionId(info.PartitionId) - , ShardId(info.ShardId) - { - if (info.KeyRange.ToBound) { - EndKeyPrefix = *info.KeyRange.ToBound; - } - } - - }; // TPartitionInfo - - TVector Schema; - TVector Partitions; - - }; // TKeyDesc - TStringBuf GetLogPrefix() const { if (!LogPrefix) { LogPrefix = TStringBuilder() @@ -453,11 +430,11 @@ class TCdcChangeSenderMain return false; } - static TVector MakePartitionIds(const TVector& partitions) { + static TVector MakePartitionIds(const TVector& partitions) { TVector result(Reserve(partitions.size())); for (const auto& partition : partitions) { - result.push_back(partition.PartitionId); + result.push_back(partition.ShardId); } return result; @@ -587,16 +564,16 @@ class TCdcChangeSenderMain const auto& pqDesc = entry.PQGroupInfo->Description; const auto& pqConfig = pqDesc.GetPQTabletConfig(); - KeyDesc = MakeHolder(); + TVector schema; PartitionToShard.clear(); - KeyDesc->Schema.reserve(pqConfig.PartitionKeySchemaSize()); + schema.reserve(pqConfig.PartitionKeySchemaSize()); for (const auto& keySchema : pqConfig.GetPartitionKeySchema()) { // TODO: support pg types - KeyDesc->Schema.push_back(NScheme::TTypeInfo(keySchema.GetTypeId())); + schema.push_back(NScheme::TTypeInfo(keySchema.GetTypeId())); } - TSet partitions(KeyDesc->Schema); + TSet partitions(schema); THashSet shards; for (const auto& partition : pqDesc.GetPartitions()) { @@ -606,8 +583,8 @@ class TCdcChangeSenderMain PartitionToShard.emplace(partitionId, shardId); auto keyRange = TPartitionKeyRange::Parse(partition.GetKeyRange()); - Y_ABORT_UNLESS(!keyRange.FromBound || keyRange.FromBound->GetCells().size() == KeyDesc->Schema.size()); - Y_ABORT_UNLESS(!keyRange.ToBound || keyRange.ToBound->GetCells().size() == KeyDesc->Schema.size()); + Y_ABORT_UNLESS(!keyRange.FromBound || keyRange.FromBound->GetCells().size() == schema.size()); + Y_ABORT_UNLESS(!keyRange.ToBound || keyRange.ToBound->GetCells().size() == schema.size()); partitions.insert({partitionId, shardId, std::move(keyRange)}); shards.insert(shardId); @@ -617,7 +594,8 @@ class TCdcChangeSenderMain bool isFirst = true; const TPQPartitionInfo* prev = nullptr; - KeyDesc->Partitions.reserve(partitions.size()); + TVector partitioning; + partitioning.reserve(partitions.size()); for (const auto& cur : partitions) { if (isFirst) { isFirst = false; @@ -629,7 +607,16 @@ class TCdcChangeSenderMain // TODO: compare cells } - KeyDesc->Partitions.emplace_back(cur); + auto& part = partitioning.emplace_back(cur.PartitionId); // TODO: double-check that it is right partitioning + + if (cur.KeyRange.ToBound) { + part.Range = NKikimr::TKeyDesc::TPartitionRangeInfo{ + .EndKeyPrefix = *cur.KeyRange.ToBound, + }; + } else { + part.Range = NKikimr::TKeyDesc::TPartitionRangeInfo{}; + } + prev = &cur; } @@ -641,7 +628,10 @@ class TCdcChangeSenderMain const bool versionChanged = !TopicVersion || TopicVersion != topicVersion; TopicVersion = topicVersion; - CreateSenders(MakePartitionIds(KeyDesc->Partitions), versionChanged); + KeyDesc = NKikimr::TKeyDesc::CreateMiniKeyDesc(schema); + KeyDesc->Partitioning = std::make_shared>(std::move(partitioning)); + + CreateSenders(MakePartitionIds(*KeyDesc->Partitioning), versionChanged); Become(&TThis::StateMain); } @@ -651,60 +641,19 @@ class TCdcChangeSenderMain return StateBase(ev); } - TActorId GetChangeServer() const override { - return DataShard.ActorId; - } - void Resolve() override { ResolveCdcStream(); } bool IsResolved() const override { - return KeyDesc && KeyDesc->Partitions; + return KeyDesc && KeyDesc->Partitioning; } - ui64 GetPartitionId(NChangeExchange::IChangeRecord::TPtr record) const override { - Y_ABORT_UNLESS(KeyDesc); - Y_ABORT_UNLESS(KeyDesc->Partitions); - - switch (Stream.Format) { - case NKikimrSchemeOp::ECdcStreamFormatProto: { - const auto range = TTableRange(record->Get()->GetKey()); - Y_ABORT_UNLESS(range.Point); - - TVector::const_iterator it = LowerBound( - KeyDesc->Partitions.begin(), KeyDesc->Partitions.end(), true, - [&](const TKeyDesc::TPartitionInfo& partition, bool) { - const int compares = CompareBorders( - partition.EndKeyPrefix.GetCells(), range.From, - partition.IsInclusive || partition.IsPoint, - range.InclusiveFrom || range.Point, KeyDesc->Schema - ); - - return (compares < 0); - } - ); - - Y_ABORT_UNLESS(it != KeyDesc->Partitions.end()); - return it->PartitionId; - } - - case NKikimrSchemeOp::ECdcStreamFormatJson: - case NKikimrSchemeOp::ECdcStreamFormatDynamoDBStreamsJson: - case NKikimrSchemeOp::ECdcStreamFormatDebeziumJson: { - using namespace NKikimr::NDataStreams::V1; - const auto hashKey = HexBytesToDecimal(record->Get()->GetPartitionKey() /* MD5 */); - return ShardFromDecimal(hashKey, KeyDesc->Partitions.size()); - } - - default: { - Y_FAIL_S("Unknown format" - << ": format# " << static_cast(Stream.Format)); - } - } - } + const TVector& GetPartitions() const override { return KeyDesc->GetPartitions(); } + const TVector& GetSchema() const override { return KeyDesc->KeyColumnTypes; } + NKikimrSchemeOp::ECdcStreamFormat GetStreamFormat() const override { return Stream.Format; } - IActor* CreateSender(ui64 partitionId) override { + IActor* CreateSender(ui64 partitionId) const override { Y_ABORT_UNLESS(PartitionToShard.contains(partitionId)); const auto shardId = PartitionToShard.at(partitionId); return new TCdcChangeSenderPartition(SelfId(), DataShard, partitionId, shardId, Stream); @@ -717,7 +666,8 @@ class TCdcChangeSenderMain void Handle(NChangeExchange::TEvChangeExchange::TEvRecords::TPtr& ev) { LOG_D("Handle " << ev->Get()->ToString()); - ProcessRecords(std::move(ev->Get()->Records)); + auto& records = std::get>>(ev->Get()->Records)->Records; + ProcessRecords(std::move(records)); } void Handle(NChangeExchange::TEvChangeExchange::TEvForgetRecords::TPtr& ev) { @@ -764,7 +714,7 @@ public: explicit TCdcChangeSenderMain(const TDataShardId& dataShard, const TPathId& streamPathId) : TActorBootstrapped() - , TBaseChangeSender(this, this, streamPathId) + , TBaseChangeSender(this, this, this, dataShard.ActorId, streamPathId) , DataShard(dataShard) , TopicVersion(0) { @@ -803,7 +753,7 @@ private: TUserTable::TCdcStream Stream; TPathId TopicPathId; ui64 TopicVersion; - THolder KeyDesc; + THolder KeyDesc; THashMap PartitionToShard; }; // TCdcChangeSenderMain diff --git a/ydb/core/tx/datashard/datashard_change_sending.cpp b/ydb/core/tx/datashard/datashard_change_sending.cpp index eb6486cae78..181f3fdf8d2 100644 --- a/ydb/core/tx/datashard/datashard_change_sending.cpp +++ b/ydb/core/tx/datashard/datashard_change_sending.cpp @@ -51,7 +51,7 @@ class TDataShard::TTxRequestChangeRecords: public TTransactionBase { struct TLoadResult { NTable::EReady Ready; - IChangeRecord::TPtr Record; + TIntrusivePtr Record; TLoadResult() = default; TLoadResult(NTable::EReady ready) @@ -59,7 +59,7 @@ class TDataShard::TTxRequestChangeRecords: public TTransactionBase { { } - explicit TLoadResult(NTable::EReady ready, IChangeRecord::TPtr record) + explicit TLoadResult(NTable::EReady ready, TIntrusivePtr record) : Ready(ready) , Record(record) { @@ -233,7 +233,7 @@ public: LOG_DEBUG_S(ctx, NKikimrServices::TX_DATASHARD, "Send " << records.size() << " change records" << ": to# " << to << ", at tablet# " << Self->TabletID()); - ctx.Send(to, new NChangeExchange::TEvChangeExchange::TEvRecords(std::move(records))); + ctx.Send(to, new NChangeExchange::TEvChangeExchange::TEvRecords(std::make_shared>(std::move(records)))); } size_t forgotten = 0; @@ -274,7 +274,7 @@ private: static constexpr size_t MemLimit = 512_KB; size_t MemUsage = 0; - THashMap> RecordsToSend; + THashMap> RecordsToSend; THashMap> RecordsToForget; }; // TTxRequestChangeRecords diff --git a/ydb/core/tx/datashard/datashard_ut_change_collector.cpp b/ydb/core/tx/datashard/datashard_ut_change_collector.cpp index 3496a46c29e..bda3c1525d4 100644 --- a/ydb/core/tx/datashard/datashard_ut_change_collector.cpp +++ b/ydb/core/tx/datashard/datashard_ut_change_collector.cpp @@ -79,7 +79,7 @@ auto GetChangeRecordsWithDetails(TTestActorRuntime& runtime, const TActorId& sen const auto details = GetChangeRecordDetails(runtime, sender, tabletId); UNIT_ASSERT_VALUES_EQUAL(records.size(), details.size()); - THashMap> result; + THashMap> result; for (size_t i = 0; i < records.size(); ++i) { const auto& record = records.at(i); const auto& detail = details.at(i); @@ -88,7 +88,7 @@ auto GetChangeRecordsWithDetails(TTestActorRuntime& runtime, const TActorId& sen const auto& pathId = std::get<4>(record); auto it = result.find(pathId); if (it == result.end()) { - it = result.emplace(pathId, TVector()).first; + it = result.emplace(pathId, TVector()).first; } it->second.push_back( @@ -359,7 +359,7 @@ Y_UNIT_TEST_SUITE(AsyncIndexChangeCollector) { UNIT_ASSERT_VALUES_EQUAL(expected.size(), actual.size()); for (size_t i = 0; i < expected.size(); ++i) { UNIT_ASSERT_VALUES_EQUAL(expected.at(i), TStructRecordBase::Parse(actual.at(i)->GetBody(), tagToName)); - UNIT_ASSERT_VALUES_EQUAL(actual.at(i)->template Get()->GetSchemaVersion(), entry.TableId.SchemaVersion); + UNIT_ASSERT_VALUES_EQUAL(actual.at(i)->GetSchemaVersion(), entry.TableId.SchemaVersion); } } } @@ -763,7 +763,7 @@ Y_UNIT_TEST_SUITE(CdcStreamChangeCollector) { UNIT_ASSERT_VALUES_EQUAL(expected.size(), actual.size()); for (size_t i = 0; i < expected.size(); ++i) { UNIT_ASSERT_VALUES_EQUAL(expected.at(i), TStructRecordBase::Parse(actual.at(i)->GetBody(), tagToName)); - UNIT_ASSERT_VALUES_EQUAL(actual.at(i)->Get()->GetSchemaVersion(), entry.TableId.SchemaVersion); + UNIT_ASSERT_VALUES_EQUAL(actual.at(i)->GetSchemaVersion(), entry.TableId.SchemaVersion); } } } diff --git a/ydb/core/tx/replication/service/json_change_record.h b/ydb/core/tx/replication/service/json_change_record.h index 89f5076fa6c..e9012b0826f 100644 --- a/ydb/core/tx/replication/service/json_change_record.h +++ b/ydb/core/tx/replication/service/json_change_record.h @@ -1,6 +1,8 @@ #pragma once +#include #include +#include #include #include #include @@ -13,6 +15,7 @@ #include #include #include +#include namespace NKikimr::NReplication::NService { @@ -48,6 +51,35 @@ public: TConstArrayRef GetKey(TMemoryPool& pool) const; TConstArrayRef GetKey() const; + ui64 ResolvePartitionId(NChangeExchange::IChangeSenderResolver* const resolver) const override { + const auto& partitions = resolver->GetPartitions(); + Y_ABORT_UNLESS(partitions); + const auto& schema = resolver->GetSchema(); + const auto streamFormat = resolver->GetStreamFormat(); + Y_ABORT_UNLESS(streamFormat == NKikimrSchemeOp::ECdcStreamFormatJson); + + // MemoryPool.Clear(); + const auto range = TTableRange(GetKey(/* MemoryPool */)); + Y_ABORT_UNLESS(range.Point); + + const auto it = LowerBound( + partitions.cbegin(), partitions.cend(), true, + [&](const auto& partition, bool) { + const int compares = CompareBorders( + partition.Range->EndKeyPrefix.GetCells(), range.From, + partition.Range->IsInclusive || partition.Range->IsPoint, + range.InclusiveFrom || range.Point, schema + ); + + return (compares < 0); + } + ); + + Y_ABORT_UNLESS(it != partitions.end()); + return it->ShardId; + } + + using TPtr = TIntrusivePtr; private: TString SourceId; NJson::TJsonValue JsonBody; @@ -82,6 +114,32 @@ public: } +namespace NKikimr { + +template <> +struct TChangeRecordContainer + : public TBaseChangeRecordContainer +{ + TChangeRecordContainer() = default; + + explicit TChangeRecordContainer(TVector&& records) + : Records(std::move(records)) + {} + + + TVector Records; + + TString Out() override { + return TStringBuilder() << "[" << JoinSeq(",", Records) << "]"; + } +}; + +} + Y_DECLARE_OUT_SPEC(inline, NKikimr::NReplication::NService::TChangeRecord, out, value) { return value.Out(out); } + +Y_DECLARE_OUT_SPEC(inline, NKikimr::NReplication::NService::TChangeRecord::TPtr, out, value) { + return value->Out(out); +} diff --git a/ydb/core/tx/replication/service/table_writer.cpp b/ydb/core/tx/replication/service/table_writer.cpp index 345c9f19c63..7e25c543e2e 100644 --- a/ydb/core/tx/replication/service/table_writer.cpp +++ b/ydb/core/tx/replication/service/table_writer.cpp @@ -77,9 +77,12 @@ class TTablePartitionWriter: public TActorBootstrapped { tableId.SetSchemaVersion(TableId.SchemaVersion); TString source; - for (auto recordPtr : ev->Get()->Records) { + + auto& records = std::get>>(ev->Get()->Records)->Records; + + for (auto recordPtr : records) { MemoryPool.Clear(); - const auto& record = *recordPtr->Get(); + const auto& record = *recordPtr; record.Serialize(*event->Record.AddChanges(), MemoryPool); if (!source) { @@ -194,8 +197,9 @@ private: class TLocalTableWriter : public TActor - , public NChangeExchange::TBaseChangeSender + , public NChangeExchange::TBaseChangeSender , public NChangeExchange::IChangeSenderResolver + , public NChangeExchange::ISenderFactory , private NSchemeCache::TSchemeCacheHelpers { TStringBuf GetLogPrefix() const { @@ -265,8 +269,8 @@ class TLocalTableWriter return result; } - TActorId GetChangeServer() const override { - return SelfId(); + void Registered(TActorSystem*, const TActorId&) override { + ChangeServer = SelfId(); } void Resolve() override { @@ -402,34 +406,13 @@ class TLocalTableWriter Resolving = false; } - IActor* CreateSender(ui64 partitionId) override { + IActor* CreateSender(ui64 partitionId) const override { return new TTablePartitionWriter(SelfId(), partitionId, TTableId(PathId, Schema->Version)); } - ui64 GetPartitionId(NChangeExchange::IChangeRecord::TPtr record) const override { - Y_ABORT_UNLESS(KeyDesc); - Y_ABORT_UNLESS(KeyDesc->GetPartitions()); - - MemoryPool.Clear(); - const auto range = TTableRange(record->Get()->GetKey(MemoryPool)); - Y_ABORT_UNLESS(range.Point); - - TVector::const_iterator it = LowerBound( - KeyDesc->GetPartitions().begin(), KeyDesc->GetPartitions().end(), true, - [&](const TKeyDesc::TPartitionInfo& partition, bool) { - const int compares = CompareBorders( - partition.Range->EndKeyPrefix.GetCells(), range.From, - partition.Range->IsInclusive || partition.Range->IsPoint, - range.InclusiveFrom || range.Point, KeyDesc->KeyColumnTypes - ); - - return (compares < 0); - } - ); - - Y_ABORT_UNLESS(it != KeyDesc->GetPartitions().end()); - return it->ShardId; - } + const TVector& GetPartitions() const override { return KeyDesc->GetPartitions(); }; + const TVector& GetSchema() const override { return KeyDesc->KeyColumnTypes; } + NKikimrSchemeOp::ECdcStreamFormat GetStreamFormat() const override { return NKikimrSchemeOp::ECdcStreamFormatJson; } void Handle(TEvWorker::TEvData::TPtr& ev) { LOG_D("Handle " << ev->Get()->ToString()); @@ -455,7 +438,7 @@ class TLocalTableWriter void Handle(NChangeExchange::TEvChangeExchange::TEvRequestRecords::TPtr& ev) { LOG_D("Handle " << ev->Get()->ToString()); - TVector records(::Reserve(ev->Get()->Records.size())); + TVector records(::Reserve(ev->Get()->Records.size())); for (const auto& record : ev->Get()->Records) { auto it = PendingRecords.find(record.Order); @@ -517,7 +500,7 @@ public: explicit TLocalTableWriter(const TPathId& tablePathId) : TActor(&TThis::StateWork) - , TBaseChangeSender(this, this, tablePathId) + , TBaseChangeSender(this, this, this, TActorId(), tablePathId) , MemoryPool(256) { } @@ -547,7 +530,7 @@ private: TLightweightSchema::TCPtr Schema; bool Resolving = false; - TMap PendingRecords; + TMap PendingRecords; }; // TLocalTableWriter -- cgit v1.3