Skip to content

Commit 82f8845

Browse files
filimonovzvonand
authored andcommitted
Merge pull request #2393 from Altinity/fix/antalya-26.6/no-json-probe-on-listing-keys
Temporary workaround: do not probe object-storage listing keys as JSON Source-PR: #2393 (#2393)
1 parent d400cad commit 82f8845

47 files changed

Lines changed: 414 additions & 390 deletions

Some content is hidden

Large Commits have some content hidden by default. Use the searchbox below for content that may be hidden.

‎src/Access/Common/AccessFlags.h‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -144,7 +144,7 @@ class AccessFlags
144144
/// The same as allColumnFlags().
145145
static AccessFlags allFlagsGrantableOnColumnLevel();
146146

147-
static constexpr size_t SIZE = 256;
147+
static constexpr size_t SIZE = 512;
148148
private:
149149
using Flags = std::bitset<SIZE>;
150150
Flags flags;

‎src/Access/Common/AccessType.h‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -150,7 +150,7 @@ ENUM_ACCESS_OBJECT(Source, APPLY_FOR_SOURCE)
150150

151151

152152
/// Represents an access type which can be granted on databases, tables, columns, etc.
153-
enum class AccessType : uint8_t
153+
enum class AccessType : uint16_t
154154
{
155155
/// Macro M should be defined as M(name, aliases, node_type, parent_group_name)
156156
/// where name is identifier with underscores (instead of spaces);

‎src/Disks/DiskObjectStorage/MetadataStorages/ContentAddressed/Backend/CasObjectStorageBackend.cpp‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -13,6 +13,7 @@
1313
#include <IO/ReadHelpers.h>
1414
#include <IO/Expect404ResponseScope.h>
1515
#include <IO/ReadSettings.h>
16+
#include <IO/WriteBufferFromFileBase.h>
1617
#include <IO/WriteBufferFromString.h>
1718
#include <IO/WriteHelpers.h>
1819
#include <IO/WriteSettings.h>

‎src/Disks/DiskObjectStorage/MetadataStorages/ContentAddressed/ContentAddressedMetadataStorage.cpp‎

Lines changed: 89 additions & 87 deletions
Original file line numberDiff line numberDiff line change
@@ -496,74 +496,75 @@ Cas::GcRoundLogger ContentAddressedMetadataStorage::makeGcRoundLogger() const
496496
auto log = ctx->getContentAddressedGarbageCollectionLog();
497497
if (!log)
498498
return;
499-
ContentAddressedGarbageCollectionLogElement e;
500-
const auto now = std::chrono::system_clock::now();
501-
e.event_time = std::chrono::system_clock::to_time_t(now);
502-
e.event_time_microseconds = timeInMicroseconds(now);
503-
switch (r.event_type)
504-
{
505-
case Cas::GcRoundLogRecord::EventType::Start:
506-
e.event_type = ContentAddressedGarbageCollectionLogElement::START;
507-
break;
508-
case Cas::GcRoundLogRecord::EventType::Finish:
509-
e.event_type = ContentAddressedGarbageCollectionLogElement::FINISH;
510-
break;
511-
case Cas::GcRoundLogRecord::EventType::Phase:
512-
e.event_type = ContentAddressedGarbageCollectionLogElement::PHASE;
513-
break;
514-
}
515-
e.disk_name = r.disk_name.empty() ? disk : r.disk_name;
516-
e.srid = r.srid;
517-
e.gc_id = r.gc_id;
518-
e.trigger = r.trigger == Cas::GcRoundLogRecord::Trigger::Manual
519-
? ContentAddressedGarbageCollectionLogElement::MANUAL
520-
: ContentAddressedGarbageCollectionLogElement::SCHEDULED;
521-
switch (r.outcome)
522-
{
523-
case Cas::GcRoundLogRecord::Outcome::Unknown:
524-
e.outcome = ContentAddressedGarbageCollectionLogElement::UNKNOWN;
525-
break;
526-
case Cas::GcRoundLogRecord::Outcome::Success:
527-
e.outcome = ContentAddressedGarbageCollectionLogElement::SUCCESS;
528-
break;
529-
case Cas::GcRoundLogRecord::Outcome::NotALeader:
530-
e.outcome = ContentAddressedGarbageCollectionLogElement::NOT_A_LEADER;
531-
break;
532-
case Cas::GcRoundLogRecord::Outcome::Failed:
533-
e.outcome = ContentAddressedGarbageCollectionLogElement::FAILED;
534-
break;
535-
case Cas::GcRoundLogRecord::Outcome::Deferred:
536-
e.outcome = ContentAddressedGarbageCollectionLogElement::DEFERRED;
537-
break;
538-
case Cas::GcRoundLogRecord::Outcome::Aborted:
539-
e.outcome = ContentAddressedGarbageCollectionLogElement::ABORTED;
540-
break;
541-
case Cas::GcRoundLogRecord::Outcome::Stopped:
542-
e.outcome = ContentAddressedGarbageCollectionLogElement::STOPPED;
543-
break;
544-
}
545-
e.round = r.round;
546-
e.candidates_marked = r.candidates_marked;
547-
e.objects_deleted = r.objects_deleted;
548-
e.objects_absent = r.objects_absent;
549-
e.objects_replaced = r.objects_replaced;
550-
e.objects_spared = r.objects_spared;
551-
e.manifests_deleted = r.manifests_deleted;
552-
e.entries_condemned = r.entries_condemned;
553-
e.entries_graduated = r.entries_graduated;
554-
e.entries_redeleted = r.entries_redeleted;
555-
e.fence_outs = r.fence_outs;
556-
e.anomalies = r.anomalies;
557-
e.duration_ms = r.duration_ms;
558-
e.error = r.error;
559-
e.error_code = r.error_code;
560-
e.profile_events = r.profile_events;
561-
e.round_id = r.round_id;
562-
e.phase = r.phase;
563-
e.phase_duration_microseconds = r.phase_duration_microseconds;
564-
e.phase_metrics = r.phase_metrics;
565499
/// Best-effort: SystemLog::add never blocks GC; a full queue drops the row with a warning.
566-
log->add(std::move(e));
500+
log->add([&](ContentAddressedGarbageCollectionLogElement & e)
501+
{
502+
const auto now = std::chrono::system_clock::now();
503+
e.event_time = std::chrono::system_clock::to_time_t(now);
504+
e.event_time_microseconds = timeInMicroseconds(now);
505+
switch (r.event_type)
506+
{
507+
case Cas::GcRoundLogRecord::EventType::Start:
508+
e.event_type = ContentAddressedGarbageCollectionLogElement::START;
509+
break;
510+
case Cas::GcRoundLogRecord::EventType::Finish:
511+
e.event_type = ContentAddressedGarbageCollectionLogElement::FINISH;
512+
break;
513+
case Cas::GcRoundLogRecord::EventType::Phase:
514+
e.event_type = ContentAddressedGarbageCollectionLogElement::PHASE;
515+
break;
516+
}
517+
e.disk_name = r.disk_name.empty() ? disk : r.disk_name;
518+
e.srid = r.srid;
519+
e.gc_id = r.gc_id;
520+
e.trigger = r.trigger == Cas::GcRoundLogRecord::Trigger::Manual
521+
? ContentAddressedGarbageCollectionLogElement::MANUAL
522+
: ContentAddressedGarbageCollectionLogElement::SCHEDULED;
523+
switch (r.outcome)
524+
{
525+
case Cas::GcRoundLogRecord::Outcome::Unknown:
526+
e.outcome = ContentAddressedGarbageCollectionLogElement::UNKNOWN;
527+
break;
528+
case Cas::GcRoundLogRecord::Outcome::Success:
529+
e.outcome = ContentAddressedGarbageCollectionLogElement::SUCCESS;
530+
break;
531+
case Cas::GcRoundLogRecord::Outcome::NotALeader:
532+
e.outcome = ContentAddressedGarbageCollectionLogElement::NOT_A_LEADER;
533+
break;
534+
case Cas::GcRoundLogRecord::Outcome::Failed:
535+
e.outcome = ContentAddressedGarbageCollectionLogElement::FAILED;
536+
break;
537+
case Cas::GcRoundLogRecord::Outcome::Deferred:
538+
e.outcome = ContentAddressedGarbageCollectionLogElement::DEFERRED;
539+
break;
540+
case Cas::GcRoundLogRecord::Outcome::Aborted:
541+
e.outcome = ContentAddressedGarbageCollectionLogElement::ABORTED;
542+
break;
543+
case Cas::GcRoundLogRecord::Outcome::Stopped:
544+
e.outcome = ContentAddressedGarbageCollectionLogElement::STOPPED;
545+
break;
546+
}
547+
e.round = r.round;
548+
e.candidates_marked = r.candidates_marked;
549+
e.objects_deleted = r.objects_deleted;
550+
e.objects_absent = r.objects_absent;
551+
e.objects_replaced = r.objects_replaced;
552+
e.objects_spared = r.objects_spared;
553+
e.manifests_deleted = r.manifests_deleted;
554+
e.entries_condemned = r.entries_condemned;
555+
e.entries_graduated = r.entries_graduated;
556+
e.entries_redeleted = r.entries_redeleted;
557+
e.fence_outs = r.fence_outs;
558+
e.anomalies = r.anomalies;
559+
e.duration_ms = r.duration_ms;
560+
e.error = r.error;
561+
e.error_code = r.error_code;
562+
e.profile_events = r.profile_events;
563+
e.round_id = r.round_id;
564+
e.phase = r.phase;
565+
e.phase_duration_microseconds = r.phase_duration_microseconds;
566+
e.phase_metrics = r.phase_metrics;
567+
});
567568
};
568569
}
569570

@@ -587,27 +588,28 @@ Cas::CasEventSink ContentAddressedMetadataStorage::makeCasEventSink() const
587588
auto log = ctx->getContentAddressedLog();
588589
if (!log)
589590
return;
590-
ContentAddressedLogElement e;
591-
const auto now = std::chrono::system_clock::now();
592-
e.event_time = std::chrono::system_clock::to_time_t(now);
593-
e.event_time_microseconds = timeInMicroseconds(now);
594-
e.event_type = toString(ev.type);
595-
e.disk_name = disk;
596-
e.namespace_ = std::move(ev.namespace_);
597-
e.ref_name = std::move(ev.ref_name);
598-
e.object_kind = toString(ev.object_kind);
599-
e.object_hash = std::move(ev.object_hash);
600-
e.token = std::move(ev.token);
601-
e.round = ev.round;
602-
e.gen = ev.gen;
603-
e.at_version = ev.at_version;
604-
e.outcome = std::move(ev.outcome);
605-
e.reason = std::move(ev.reason);
606-
e.thread_id = getThreadId();
607-
e.query_id = CurrentThread::getQueryId();
608-
e.detail = std::move(ev.detail);
609591
/// Best-effort: SystemLog::add never blocks the Core; a full queue drops the row with a warning.
610-
log->add(std::move(e));
592+
log->add([&](ContentAddressedLogElement & e)
593+
{
594+
const auto now = std::chrono::system_clock::now();
595+
e.event_time = std::chrono::system_clock::to_time_t(now);
596+
e.event_time_microseconds = timeInMicroseconds(now);
597+
e.event_type = toString(ev.type);
598+
e.disk_name = disk;
599+
e.namespace_ = std::move(ev.namespace_);
600+
e.ref_name = std::move(ev.ref_name);
601+
e.object_kind = toString(ev.object_kind);
602+
e.object_hash = std::move(ev.object_hash);
603+
e.token = std::move(ev.token);
604+
e.round = ev.round;
605+
e.gen = ev.gen;
606+
e.at_version = ev.at_version;
607+
e.outcome = std::move(ev.outcome);
608+
e.reason = std::move(ev.reason);
609+
e.thread_id = getThreadId();
610+
e.query_id = CurrentThread::getQueryId();
611+
e.detail = std::move(ev.detail);
612+
});
611613
};
612614
}
613615

‎src/Disks/DiskObjectStorage/MetadataStorages/ContentAddressed/ContentAddressedTransaction.cpp‎

Lines changed: 14 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -8,7 +8,7 @@
88
#include <IO/ReadBufferFromMemory.h>
99
#include <IO/WriteBufferFromFile.h>
1010
#include <IO/copyData.h>
11-
#include <Disks/IO/WriteBufferWithFinalizeCallback.h>
11+
#include <Disks/IO/WriteBufferInlineOrBlob.h>
1212
#include <Disks/IDiskTransaction.h>
1313
#include <Common/thread_local_rng.h>
1414
#include <Common/config_version.h>
@@ -796,10 +796,10 @@ std::unique_ptr<WriteBufferFromFileBase> ContentAddressedTransaction::tryCreateW
796796
throw Exception(ErrorCodes::NOT_IMPLEMENTED,
797797
"Autocommit writes are not supported for content part files on a content-addressed disk");
798798

799-
auto inner = writeFile(path, buf_size, mode, settings);
800-
auto commit_callback = [owner](size_t) mutable { owner->commit(); };
801-
return std::make_unique<WriteBufferWithFinalizeCallback>(
802-
std::move(inner), std::move(commit_callback), path, /*create_blob_if_empty=*/true);
799+
auto create_inner = [this, path, buf_size, mode, settings] { return writeFile(path, buf_size, mode, settings); };
800+
auto commit_callback = [owner](FinalizeResult) mutable { owner->commit(); };
801+
return std::make_unique<WriteBufferInlineOrBlob>(
802+
path, /*max_inline_bytes=*/0, /*create_blob_if_empty=*/true, std::move(create_inner), std::move(commit_callback), buf_size);
803803
}
804804

805805
/// Non-autocommit (or verbatim autocommit): pin the owning disk transaction for the returned buffer's
@@ -809,10 +809,10 @@ std::unique_ptr<WriteBufferFromFileBase> ContentAddressedTransaction::tryCreateW
809809
/// `owner` (which owns this ContentAddressedTransaction by shared_ptr) keeps that `this` valid until the
810810
/// buffer — and so this callback — is destroyed after finalize (the lifetime guarantee now
811811
/// expressed generically via `owner`). No cycle: the transaction does not hold the buffer.
812-
auto inner = writeFile(path, buf_size, mode, settings);
813-
auto keep_alive_callback = [owner](size_t) mutable {};
814-
return std::make_unique<WriteBufferWithFinalizeCallback>(
815-
std::move(inner), std::move(keep_alive_callback), path, /*create_blob_if_empty=*/true);
812+
auto create_inner = [this, path, buf_size, mode, settings] { return writeFile(path, buf_size, mode, settings); };
813+
auto keep_alive_callback = [owner](FinalizeResult) mutable {};
814+
return std::make_unique<WriteBufferInlineOrBlob>(
815+
path, /*max_inline_bytes=*/0, /*create_blob_if_empty=*/true, std::move(create_inner), std::move(keep_alive_callback), buf_size);
816816
}
817817

818818
std::unique_ptr<WriteBufferFromFileBase> ContentAddressedTransaction::writeFile(
@@ -1831,10 +1831,10 @@ CaContentWriteBuffer::CaContentWriteBuffer(
18311831
std::string temp_dir,
18321832
Cas::BlobHashAlgo hash_algo,
18331833
size_t buf_size,
1834-
bool use_adaptive_buffer_size,
1834+
bool use_adaptive_buffer_size_,
18351835
size_t adaptive_buffer_initial_size,
18361836
OnFinalized on_finalized_)
1837-
: WriteBufferFromFileBase(clampCasWriteBufferSize(use_adaptive_buffer_size ? adaptive_buffer_initial_size : buf_size), nullptr, 0)
1837+
: WriteBufferFromFileBase(clampCasWriteBufferSize(use_adaptive_buffer_size_ ? adaptive_buffer_initial_size : buf_size), nullptr, 0)
18381838
, on_finalized(std::move(on_finalized_))
18391839
{
18401840
fs::create_directories(temp_dir);
@@ -1850,7 +1850,7 @@ CaContentWriteBuffer::CaContentWriteBuffer(
18501850
/*mode=*/0666,
18511851
/*existing_memory=*/nullptr,
18521852
/*alignment=*/0,
1853-
use_adaptive_buffer_size,
1853+
use_adaptive_buffer_size_,
18541854
clampCasWriteBufferSize(adaptive_buffer_initial_size));
18551855
hashing = Cas::makeBlobHashingWriteBuffer(hash_algo, *sink);
18561856
}
@@ -1861,11 +1861,11 @@ CaContentWriteBuffer::CaContentWriteBuffer(
18611861
std::string envelope_header,
18621862
Cas::BlobHashAlgo hash_algo,
18631863
size_t buf_size,
1864-
bool use_adaptive_buffer_size,
1864+
bool use_adaptive_buffer_size_,
18651865
size_t adaptive_buffer_initial_size,
18661866
OnFinalized on_finalized_,
18671867
std::function<void()> check_fence_before_finalize_)
1868-
: WriteBufferFromFileBase(clampCasWriteBufferSize(use_adaptive_buffer_size ? adaptive_buffer_initial_size : buf_size), nullptr, 0)
1868+
: WriteBufferFromFileBase(clampCasWriteBufferSize(use_adaptive_buffer_size_ ?adaptive_buffer_initial_size : buf_size), nullptr, 0)
18691869
, on_finalized(std::move(on_finalized_))
18701870
, temp_path(std::move(object_key))
18711871
, is_s3_staging(true)

‎src/Disks/DiskObjectStorage/MetadataStorages/ContentAddressed/Pool/CasMountRuntime.cpp‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,7 @@
11
#include <Disks/DiskObjectStorage/MetadataStorages/ContentAddressed/Pool/CasMountRuntime.h>
22
#include <Disks/DiskObjectStorage/MetadataStorages/ContentAddressed/Pool/CasPartWriteTxn.h>
33
#include <Common/Exception.h>
4+
#include <Common/ProfileEvents.h>
45
#include <Common/logger_useful.h>
56
#include <Common/setThreadName.h>
67
#include <Common/thread_local_rng.h>

‎src/Disks/DiskObjectStorage/ObjectStorages/IObjectStorage.cpp‎

Lines changed: 20 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -162,6 +162,26 @@ RelativePathWithMetadata::RelativePathWithMetadata(const DataFileInfo & info, st
162162

163163
RelativePathWithMetadata::CommandInTaskResponse::CommandInTaskResponse(const std::string & task)
164164
{
165+
/// TEMPORARY WORKAROUND. This constructor runs for every string passed to `RelativePathWithMetadata`,
166+
/// which includes every key returned by a `listObjects` / `iterate` call of every object storage, not only
167+
/// the task-distributor answers it exists for (`{"retry_after_us": N}` from
168+
/// `StorageObjectStorageStableTaskDistributor`). Parsing an ordinary object key as JSON throws and catches
169+
/// one `JSONException` per listed key. Besides the cost, under ASan the fake stack frame of `parseImpl`
170+
/// that exits by exception is never released (it is re-entered at the same stack depth, and `FakeStack::GC`
171+
/// frees only frames strictly below the next allocation), so every listed key leaks one frame per thread;
172+
/// once the size class is full every `__asan_stack_malloc_1` scans all 8192 slots and every small function
173+
/// on that thread becomes ~100x slower. See https://github.com/Altinity/ClickHouse/issues/2362.
174+
///
175+
/// Only try to parse strings that can be a JSON object. Object keys never start with `{`; the distributor
176+
/// answer always does. The proper fix is to stop multiplexing the command into the path field (a separate
177+
/// `ObjectInfo` kind or the versioned cluster-function protocol, see
178+
/// https://github.com/Altinity/ClickHouse/pull/1360), after which this probe goes away entirely.
179+
{
180+
const auto first = task.find_first_not_of(" \t\r\n");
181+
if (first == std::string::npos || task[first] != '{')
182+
return;
183+
}
184+
165185
Poco::JSON::Parser parser;
166186
try
167187
{

0 commit comments

Comments
 (0)