diff --git a/be/src/cloud/cloud_delta_writer.cpp b/be/src/cloud/cloud_delta_writer.cpp index 7a7fd00bb687f3..b85e83be20c0bd 100644 --- a/be/src/cloud/cloud_delta_writer.cpp +++ b/be/src/cloud/cloud_delta_writer.cpp @@ -25,6 +25,8 @@ #include "load/memtable/memtable_memory_limiter.h" #include "runtime/exec_env.h" #include "runtime/thread_context.h" +#include "storage/adaptive_thread_pool_controller.h" +#include "util/threadpool.h" namespace doris { @@ -86,8 +88,18 @@ Status CloudDeltaWriter::write(const Block* block, const DorisVector& CHECK(_is_init || _is_cancelled); { SCOPED_TIMER(_wait_flush_limit_timer); - while (_memtable_writer->flush_running_count() >= - config::memtable_flush_running_count_limit) { + auto* s3_file_upload_pool = rowset_builder()->is_s3_storage() + ? ExecEnv::GetInstance()->s3_file_upload_thread_pool() + : nullptr; + const auto need_backpressure = [this, s3_file_upload_pool] { + if (s3_file_upload_pool != nullptr) { + return s3_file_upload_pool->get_queue_size() > + AdaptiveThreadPoolController::kS3QueueBusyThreshold; + } + return _memtable_writer->flush_running_count() >= + config::memtable_flush_running_count_limit; + }; + while (need_backpressure()) { std::this_thread::sleep_for(std::chrono::milliseconds(10)); } } diff --git a/be/src/cloud/cloud_rowset_builder.cpp b/be/src/cloud/cloud_rowset_builder.cpp index 9762f2b62b6a5c..e03d347d94fe5e 100644 --- a/be/src/cloud/cloud_rowset_builder.cpp +++ b/be/src/cloud/cloud_rowset_builder.cpp @@ -21,6 +21,7 @@ #include "cloud/cloud_storage_engine.h" #include "cloud/cloud_tablet.h" #include "cloud/cloud_tablet_mgr.h" +#include "io/fs/file_system.h" #include "storage/storage_policy.h" namespace doris { @@ -135,6 +136,13 @@ const RowsetMetaSharedPtr& CloudRowsetBuilder::rowset_meta() { return _rowset_writer->rowset_meta(); } +bool CloudRowsetBuilder::is_s3_storage() const { + if (_rowset_writer == nullptr) { + return false; + } + return _rowset_writer->context().fs()->type() == io::FileSystemType::S3; +} + Status CloudRowsetBuilder::set_txn_related_info() { if (_tablet->enable_unique_key_merge_on_write()) { // For empty rowsets when skip_writing_empty_rowset_metadata=true, diff --git a/be/src/cloud/cloud_rowset_builder.h b/be/src/cloud/cloud_rowset_builder.h index cec8cfed979857..5549a2616591d9 100644 --- a/be/src/cloud/cloud_rowset_builder.h +++ b/be/src/cloud/cloud_rowset_builder.h @@ -37,6 +37,8 @@ class CloudRowsetBuilder final : public BaseRowsetBuilder { const RowsetMetaSharedPtr& rowset_meta(); + bool is_s3_storage() const; + Status set_txn_related_info(); void set_skip_writing_rowset_metadata(bool skip) { _skip_writing_rowset_metadata = skip; }