summaryrefslogtreecommitdiffstats
diff options
context:
space:
mode:
authorapollo1321 <[email protected]>2026-07-22 11:41:02 +0300
committerapollo1321 <[email protected]>2026-07-22 12:28:54 +0300
commita08cbf0881a0c0bbd5d9fa98f149ddefa73fa5e6 (patch)
tree8fc78fb11d1a7bc8a4a4f5d4cba93be132f89899
parenta4e1acc61d1f1b147aaa50b173cb9f78a503a5f9 (diff)
YT-28707: Support pipelining in journal chunk writer
Better disks utilization (\~85% -\> \~98%) commit_hash:9f87512a891cba9ef921ed1cd910199b562b0695
-rw-r--r--yt/yt/client/api/config.cpp4
-rw-r--r--yt/yt/client/api/config.h3
-rw-r--r--yt/yt/client/api/journal_client.cpp1
-rw-r--r--yt/yt/client/api/journal_client.h2
4 files changed, 10 insertions, 0 deletions
diff --git a/yt/yt/client/api/config.cpp b/yt/yt/client/api/config.cpp
index b67e6ddcccf..885442301c4 100644
--- a/yt/yt/client/api/config.cpp
+++ b/yt/yt/client/api/config.cpp
@@ -100,6 +100,10 @@ void TJournalChunkWriterConfig::Register(TRegistrar registrar)
.Default(100'000);
registrar.Parameter("max_flush_data_size", &TThis::MaxFlushDataSize)
.Default(100_MB);
+ registrar.Parameter("max_in_flight_flush_count", &TThis::MaxInFlightFlushCount)
+ .Default(1)
+ .GreaterThanOrEqual(1)
+ .DontSerializeDefault();
registrar.Parameter("prefer_local_host", &TThis::PreferLocalHost)
.Default(true);
diff --git a/yt/yt/client/api/config.h b/yt/yt/client/api/config.h
index ef39a4ba3ce..eb16e7cc364 100644
--- a/yt/yt/client/api/config.h
+++ b/yt/yt/client/api/config.h
@@ -172,6 +172,9 @@ struct TJournalChunkWriterConfig
int MaxFlushRowCount;
i64 MaxFlushDataSize;
+ //! Maximum number of inflight PutBlocks/Flush requests per replica.
+ int MaxInFlightFlushCount;
+
bool PreferLocalHost;
TDuration NodeRpcTimeout;
diff --git a/yt/yt/client/api/journal_client.cpp b/yt/yt/client/api/journal_client.cpp
index 96854d1174e..cdc985a6b5f 100644
--- a/yt/yt/client/api/journal_client.cpp
+++ b/yt/yt/client/api/journal_client.cpp
@@ -30,6 +30,7 @@ TJournalWriterPerformanceCounters::TJournalWriterPerformanceCounters(const NProf
MediumWrittenBytes = profiler.Counter("/medium_written_bytes");
JournalWrittenBytes = profiler.Counter("/journal_written_bytes");
IORequestCount = profiler.Counter("/io_request_count");
+ PipelinedFlushCount = profiler.Counter("/pipelined_flush_count");
}
////////////////////////////////////////////////////////////////////////////////
diff --git a/yt/yt/client/api/journal_client.h b/yt/yt/client/api/journal_client.h
index 51162a9dbd6..c0679a4c1ac 100644
--- a/yt/yt/client/api/journal_client.h
+++ b/yt/yt/client/api/journal_client.h
@@ -54,6 +54,8 @@ struct TJournalWriterPerformanceCounters
NProfiling::TCounter IORequestCount;
NProfiling::TCounter JournalWrittenBytes;
+ NProfiling::TCounter PipelinedFlushCount;
+
IJournalWritesObserverPtr JournalWritesObserver;
};