Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions src/Common/FailPoint.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -35,6 +35,8 @@ static struct InitFiu
ONCE(replicated_merge_tree_commit_zk_fail_after_op) \
ONCE(replicated_queue_fail_next_entry) \
REGULAR(replicated_queue_unfail_entries) \
ONCE(transaction_metadata_store_fail) \
ONCE(transaction_mutation_csn_store_fail) \
ONCE(replicated_merge_tree_insert_quorum_fail_0) \
REGULAR(replicated_merge_tree_commit_zk_fail_when_recovering_from_hw_fault) \
REGULAR(rmt_dedup_conflict_part_name_missing) \
Expand Down
106 changes: 92 additions & 14 deletions src/Interpreters/MergeTreeTransaction.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -9,7 +9,9 @@
#include <Common/FailPoint.h>
#include <Common/ThreadPool.h>
#include <Common/TransactionID.h>
#include <Common/Stopwatch.h>
#include <Common/ZooKeeper/IKeeper.h>
#include <Common/logger_useful.h>
#include <Common/noexcept_scope.h>

#include <base/sleep.h>
Expand All @@ -31,6 +33,77 @@ namespace FailPoints
extern const char transaction_after_commit_pause[];
}

namespace
{

/// A metadata write made after the commit point, or during rollback, has no one to report
/// an error to: the caller is a noexcept callback and the transaction's fate is already
/// decided in the transaction log. Instead of letting the exception terminate the server,
/// the write is retried for a bounded time. `LOGICAL_ERROR` and `NOT_IMPLEMENTED` are invariant
/// violations and are rethrown at once. When the budget is exhausted, or the server is
/// shutting down, the error is rethrown too: this keeps the old behaviour rather than
/// hiding a write that did not happen. The budget is per object; a write that hangs

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The budget is per object

The retry budget is per object, so with flapping storage errors a commit with N parts can take up to N*60 s. This is longer than receive_timeout, and KILL QUERY or max_execution_time cannot stop it. On the unknown-state path
this blocks runUpdatingThread, so every COMMIT on the server waits in waitForCSNLoaded.

I suggest using one deadline for the whole afterCommit / rollback, and still rethrow when it expires.

/// inside the storage is not interrupted.
constexpr UInt64 TRANSACTION_METADATA_STORE_RETRY_TIMEOUT_SECONDS = 60;
constexpr UInt64 TRANSACTION_METADATA_STORE_RETRY_BACKOFF_MS = 100;
constexpr UInt64 TRANSACTION_METADATA_STORE_RETRY_MAX_BACKOFF_MS = 2000;

/// `describe` is called only when a log line is written, so a callback that never fails
/// formats nothing.
template <typename Describe, typename F>
void retryMetadataStore(LoggerPtr log, Describe && describe, F && store)
{
Stopwatch watch;
UInt64 backoff_ms = TRANSACTION_METADATA_STORE_RETRY_BACKOFF_MS;
size_t attempts = 0;
while (true)
{
++attempts;
try
{
store();
if (attempts > 1)
LOG_INFO(log, "Stored transaction metadata for {} after {} attempts", describe(), attempts);
return;
}
catch (...)
{
int code = getCurrentExceptionCode();
if (code == ErrorCodes::LOGICAL_ERROR || code == ErrorCodes::NOT_IMPLEMENTED)
throw;

bool give_up = watch.elapsedSeconds() >= TRANSACTION_METADATA_STORE_RETRY_TIMEOUT_SECONDS
|| TransactionLog::instance().isShuttingDown();
if (give_up)
{
LOG_ERROR(log, "Cannot store transaction metadata for {} after {} attempts in {:.1f} s, giving up: {}",
describe(), attempts, watch.elapsedSeconds(), getCurrentExceptionMessage(false));
throw;
}

if (attempts == 1)
LOG_WARNING(log, "Cannot store transaction metadata for {}, will retry: {}", describe(), getCurrentExceptionMessage(false));
else
LOG_DEBUG(log, "Cannot store transaction metadata for {}, attempt {}: {}", describe(), attempts, getCurrentExceptionMessage(false));
}

sleepForMilliseconds(backoff_ms);
backoff_ms = std::min(backoff_ms * 2, TRANSACTION_METADATA_STORE_RETRY_MAX_BACKOFF_MS);
}
}

String partDescription(const IMergeTreeDataPart & part)
{
return fmt::format("part {} of {}", part.name, part.storage.getStorageID().getNameForLogs());
}

String mutationDescription(const IStorage & storage, const String & mutation_id)
{
return fmt::format("mutation {} of {}", mutation_id, storage.getStorageID().getNameForLogs());
}

}

static void checkNotOrdinaryDatabase(const StoragePtr & storage)
{
if (storage->getStorageID().uuid != UUIDHelpers::Nil)
Expand Down Expand Up @@ -275,6 +348,8 @@ void MergeTreeTransaction::afterCommit(CSN assigned_csn) noexcept
committed_mutations = mutations;
}

auto log = getLogger("MergeTreeTransaction");

/// Persist per-part version metadata BEFORE flipping `csn` below.
/// `csn.exchange(assigned_csn)` is the signal that `MergeTreeTransaction::waitStateChange`
/// blocks on; doing the disk-backed `setAndStore...CSN` calls first ensures that once a
Expand All @@ -288,23 +363,20 @@ void MergeTreeTransaction::afterCommit(CSN assigned_csn) noexcept
/// `setAndStore...CSN` did not complete; `TransactionLog::getCSN(tid)` returns the right
/// answer after restart.
for (const auto & part : created_parts)
{
part->version->setAndStoreCreationCSN(assigned_csn);
}
retryMetadataStore(log, [&] { return partDescription(*part); }, [&] { part->version->setAndStoreCreationCSN(assigned_csn); });

for (const auto & part : removed_parts)
{
part->version->setAndStoreRemovalCSN(assigned_csn);
}
retryMetadataStore(log, [&] { return partDescription(*part); }, [&] { part->version->setAndStoreRemovalCSN(assigned_csn); });

for (const auto & storage_and_mutation : committed_mutations)
storage_and_mutation.first->setMutationCSN(storage_and_mutation.second, assigned_csn);
retryMetadataStore(log, [&] { return mutationDescription(*storage_and_mutation.first, storage_and_mutation.second); },
[&] { storage_and_mutation.first->setMutationCSN(storage_and_mutation.second, assigned_csn); });

/// Test-only pause point. With this failpoint enabled, a regression test can verify that
/// `waitStateChange` does not return until every part has its new CSN persisted (above).
/// Not wrapped in try/catch: `pauseFailPoint` only takes a mutex and a condvar, and the
/// surrounding `setAndStore...CSN` calls already trust their callees not to throw under
/// the same `noexcept` contract.
/// Not wrapped in try/catch: `pauseFailPoint` only takes a mutex and a condvar. The writes
/// above go through retryMetadataStore, which absorbs recoverable storage errors within its
/// retry budget.
FailPointInjection::pauseFailPoint(FailPoints::transaction_after_commit_pause);

/// Flip the atomic last so that `waitStateChange` only wakes up after all metadata is durable.
Expand Down Expand Up @@ -337,9 +409,15 @@ bool MergeTreeTransaction::rollback() noexcept
parts_to_activate = removing_parts;
}

/// Forcefully stop related mutations if any
auto log = getLogger("MergeTreeTransaction");

/// Forcefully stop related mutations if any. killMutation erases the mutation from the
/// table's map before deleting its file, so a repeated call after a failed deletion is a
/// no-op; a file left behind is removed at the next load, because its transaction has
/// no CSN.
for (const auto & table_and_mutation : mutations_to_kill)
table_and_mutation.first->killMutation(table_and_mutation.second);
retryMetadataStore(log, [&] { return mutationDescription(*table_and_mutation.first, table_and_mutation.second); },

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Example: StorageMergeTree::killMutation

Do we need to check the result?
We could finish current_mutations_by_version.erase(it); but throw before removeFile. So retring doesn't help o remove the file.

Maybe we don't need retry here?

[&] { table_and_mutation.first->killMutation(table_and_mutation.second); });

/// Discard changes in active parts set
/// Remove parts that were created, restore parts that were removed (except parts that were created by this transaction too)
Expand All @@ -348,7 +426,7 @@ bool MergeTreeTransaction::rollback() noexcept
for (const auto & part : parts_to_remove)
{
/// Write special RolledBackCSN, so we will be able to cleanup transaction log
part->version->setAndStoreCreationCSN(Tx::RolledBackCSN);
retryMetadataStore(log, [&] { return partDescription(*part); }, [&] { part->version->setAndStoreCreationCSN(Tx::RolledBackCSN); });
}

for (const auto & part : parts_to_remove)
Expand All @@ -367,7 +445,7 @@ bool MergeTreeTransaction::rollback() noexcept
{
/// Clear removal_tid from version metadata file, so we will not need to distinguish TIDs that were not committed
/// and TIDs that were committed long time ago and were removed from the log on log cleanup.
part->version->setAndStoreRemovalTID(Tx::EmptyTID);
retryMetadataStore(log, [&] { return partDescription(*part); }, [&] { part->version->setAndStoreRemovalTID(Tx::EmptyTID); });
part->version->unlockRemovalTID(tid, TransactionInfoContext{part->storage.getStorageID(), part->name});
}

Expand Down
13 changes: 13 additions & 0 deletions src/Interpreters/MergeTreeTransaction/VersionMetadataOnDisk.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@
#include <Storages/MergeTree/MergeTreeData.h>
#include <Storages/MergeTree/MergeTreeSettings.h>
#include <Common/Exception.h>
#include <Common/FailPoint.h>
#include <Common/TransactionID.h>
#include <Common/logger_useful.h>
#include <base/scope_guard.h>
Expand All @@ -28,6 +29,12 @@ namespace ErrorCodes
extern const int LOGICAL_ERROR;
extern const int CANNOT_OPEN_FILE;
extern const int NOT_IMPLEMENTED;
extern const int FAULT_INJECTED;
}

namespace FailPoints
{
extern const char transaction_metadata_store_fail[];
}

namespace MergeTreeSetting
Expand Down Expand Up @@ -323,6 +330,12 @@ void VersionMetadataOnDisk::removeTmpMetadataFile()
void VersionMetadataOnDisk::storeInfoToDataPartStorage(
const MergeTreeData & storage, IDataPartStorage & data_part_storage, const VersionInfo & new_info)
{
/// Fault injection for tests: fail before any I/O, so the old file stays intact.
fiu_do_on(FailPoints::transaction_metadata_store_fail,
{
throw Exception(ErrorCodes::FAULT_INJECTED, "Injected failure while storing version metadata");
});

static constexpr auto filename = TXN_VERSION_METADATA_FILE_NAME;
static constexpr auto tmp_filename = TMP_TXN_VERSION_METADATA_FILE_NAME;

Expand Down
57 changes: 42 additions & 15 deletions src/Storages/MergeTree/MergeTreeMutationEntry.cpp
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
#include <Storages/MergeTree/MergeTreeMutationEntry.h>
#include <Common/FailPoint.h>
#include <Common/logger_useful.h>
#include <IO/Operators.h>
#include <IO/ReadHelpers.h>
Expand All @@ -17,6 +18,12 @@ namespace DB
namespace ErrorCodes
{
extern const int BAD_ARGUMENTS;
extern const int FAULT_INJECTED;
}

namespace FailPoints
{
extern const char transaction_mutation_csn_store_fail[];
}

String MergeTreeMutationEntry::versionToFileName(UInt64 block_number_)
Expand Down Expand Up @@ -60,22 +67,11 @@ MergeTreeMutationEntry::MergeTreeMutationEntry(MutationCommands commands_, DiskP
{
try
{
auto out = disk->writeFile(std::filesystem::path(path_prefix) / file_name, DBMS_DEFAULT_BUFFER_SIZE, WriteMode::Rewrite, settings);
*out << "format version: 1\n"
<< "create time: " << LocalDateTime(create_time, DateLUT::serverTimezoneInstance()) << "\n";
*out << "commands: ";
commands->writeText(*out, /* with_pure_metadata_commands = */ false);
*out << "\n";
if (tid.isNonTransactional())
{
csn = Tx::NonTransactionalCSN;
}
else
{
*out << "tid: ";
TransactionID::write(tid, *out);
*out << "\n";
}

auto out = disk->writeFile(std::filesystem::path(path_prefix) / file_name, DBMS_DEFAULT_BUFFER_SIZE, WriteMode::Rewrite, settings);
writeRecord(*out);
out->finalize();
out->sync();
}
Expand All @@ -86,6 +82,21 @@ MergeTreeMutationEntry::MergeTreeMutationEntry(MutationCommands commands_, DiskP
}
}

void MergeTreeMutationEntry::writeRecord(WriteBuffer & out) const
{
out << "format version: 1\n"
<< "create time: " << LocalDateTime(create_time, DateLUT::serverTimezoneInstance()) << "\n";
out << "commands: ";
commands->writeText(out, /* with_pure_metadata_commands = */ false);
out << "\n";
if (!tid.isNonTransactional())
{
out << "tid: ";
TransactionID::write(tid, out);
out << "\n";
}
}

void MergeTreeMutationEntry::commit(UInt64 block_number_)
{
chassert(block_number_);
Expand All @@ -111,9 +122,25 @@ void MergeTreeMutationEntry::removeFile()
void MergeTreeMutationEntry::writeCSN(CSN csn_)
{
csn = csn_;

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

nit: set csn after file replacing.

auto out = disk->writeFile(path_prefix + file_name, 256, WriteMode::Append);
/// Fault injection for tests: fail before any I/O, so the old file stays intact.
fiu_do_on(FailPoints::transaction_mutation_csn_store_fail,
{
throw Exception(ErrorCodes::FAULT_INJECTED, "Injected failure while storing mutation CSN");
});

/// The whole record is rewritten through a temporary file instead of appending the
/// `csn:` line: a write that fails half-way, or is repeated, must not leave a partial
/// or duplicated line behind, because the loader accepts exactly one.
/// The name must not collide with the constructor's `tmp_mutation_<N>.txt`, whose
/// number comes from a different counter; the `tmp_mutation_` prefix keeps it covered
/// by the startup cleanup of leftover temporary files.
String tmp_file_name = "tmp_mutation_csn_" + toString(block_number) + ".txt";
auto out = disk->writeFile(path_prefix + tmp_file_name, DBMS_DEFAULT_BUFFER_SIZE, WriteMode::Rewrite);
writeRecord(*out);
*out << "csn: " << csn << "\n";
out->finalize();
out->sync();
disk->replaceFile(path_prefix + tmp_file_name, path_prefix + file_name);
}

MergeTreeMutationEntry::MergeTreeMutationEntry(DiskPtr disk_, const String & path_prefix_, const String & file_name_)
Expand Down
5 changes: 5 additions & 0 deletions src/Storages/MergeTree/MergeTreeMutationEntry.h
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@
namespace DB
{
class IBackupEntry;
class WriteBuffer;

/// A mutation entry for non-replicated MergeTree storage engines.
/// Stores information about mutation in file mutation_*.txt.
Expand Down Expand Up @@ -79,6 +80,10 @@ struct MergeTreeMutationEntry
MergeTreeMutationEntry(DiskPtr disk_, const String & path_prefix_, const String & file_name_);

~MergeTreeMutationEntry();

private:
/// Serializes everything except the `csn:` line, in the format the loading constructor reads.
void writeRecord(WriteBuffer & out) const;
};

}
12 changes: 10 additions & 2 deletions src/Storages/StorageMergeTree.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -1128,7 +1128,13 @@ void StorageMergeTree::setMutationCSN(const String & mutation_id, CSN csn)
std::lock_guard lock(currently_processing_in_background_mutex);
auto it = current_mutations_by_version.find(version);
if (it == current_mutations_by_version.end())
throw Exception(ErrorCodes::LOGICAL_ERROR, "Cannot find mutation {}", mutation_id);
{
/// `KILL MUTATION` erases the entry before the committing transaction stores the CSN,
/// and cannot roll that transaction back any more. The parts are already mutated; the
/// file, if its deletion is still pending or failed, is resolved at the next load.
LOG_WARNING(log, "Mutation {} was killed before its CSN {} could be stored", mutation_id, csn);
return;
}
it->second.writeCSN(csn);
}

Expand Down Expand Up @@ -1567,7 +1573,9 @@ void StorageMergeTree::loadMutations()
}
else if (startsWith(it->name(), "tmp_mutation_"))
{
disk->removeFile(it->path());
/// A CSN write of an entry loaded earlier in this pass may have consumed a

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

loadMutations now calls writeCSN, which does replaceFile on mutation_N.txt, while iterateDirectory is still running. POSIX does not say whether readdir returns entries that were added or removed during iteration. The PR handles the temporary file (removeFileIfExists). Should we also handle the case where the iterator returns the recreated mutation_N.txt a second time?

/// temporary file the iterator still lists, so a missing file is not an error.
disk->removeFileIfExists(it->path());
}
}
}
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,16 @@
--- A: commit, part metadata
retried objects 1 stored after retry 1 same object 1 attempts [2]
creation_csn persisted: 1
rows after commit 1
rows after reattach 1
--- B: commit, mutation CSN
retried objects 1 stored after retry 1 same object 1 attempts [2]
csn lines in mutation file: 1
last line is the csn: 1
sum after commit 9
sum after reattach 9
mutations after reattach 1
--- C: rollback, part metadata
retried objects 1 stored after retry 1 same object 1 attempts [2]
rows after rollback 2
rows after reattach 2
Loading
Loading