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
8 changes: 8 additions & 0 deletions src/include/duckdb/storage/write_ahead_log.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,7 @@
#include "duckdb/common/enums/wal_type.hpp"
#include "duckdb/common/serializer/buffered_file_writer.hpp"
#include "duckdb/common/types/data_chunk.hpp"
#include "duckdb/common/unordered_set.hpp"
#include "duckdb/storage/block.hpp"

namespace duckdb {
Expand Down Expand Up @@ -125,13 +126,20 @@ class WriteAheadLog {
void IncrementWALEntriesCount();
void WriteCheckpoint(MetaBlockPointer meta_block);

//! Drop the pending blocks of a reverted commit
void RemovePendingCheckpointBlocks(const unordered_set<block_id_t> &block_ids);
//! Mark the pending blocks as checkpointed, called after the concurrent checkpoint has written its header
void MarkPendingBlocksAsCheckpointed();

protected:
StorageManager &storage_manager;
mutex wal_lock;
unique_ptr<BufferedFileWriter> writer;
string wal_path;
atomic<WALInitState> init_state;
optional_idx checkpoint_iteration;
//! Row group blocks written to this checkpoint WAL, they stay newly used until the running checkpoint is done
vector<block_id_t> pending_checkpoint_blocks;
};

} // namespace duckdb
5 changes: 5 additions & 0 deletions src/storage/storage_manager.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -308,6 +308,8 @@ bool StorageManager::WALStartCheckpoint(MetaBlockPointer meta_block, CheckpointO

void StorageManager::WALFinishCheckpoint(unique_lock<mutex> &) {
D_ASSERT(wal.get());
// the header is written - the next checkpoint includes the commits of the checkpoint WAL
wal->MarkPendingBlocksAsCheckpointed();

// "wal" points to the checkpoint WAL
// first check if the checkpoint WAL has been written to
Expand Down Expand Up @@ -672,15 +674,18 @@ void SingleFileStorageCommitState::RevertCommit() {
wal.Truncate(initial_wal_size);
}
auto &block_manager = storage.GetBlockManager();
unordered_set<block_id_t> reverted_blocks;
for (auto &entry : optimistically_written_data) {
for (auto &rg_entry : entry.second) {
if (rg_entry.second.row_group_data) {
for (auto &block_id : rg_entry.second.row_group_data->GetBlockIds()) {
block_manager.MarkBlockAsModified(block_id);
reverted_blocks.insert(block_id);
}
}
}
}
wal.RemovePendingCheckpointBlocks(reverted_blocks);
state = WALCommitState::TRUNCATED;
}

Expand Down
27 changes: 27 additions & 0 deletions src/storage/write_ahead_log.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,8 @@
#include "duckdb/storage/wal_entry.hpp"
#include "duckdb/main/attached_database.hpp"

#include <algorithm>

namespace duckdb {

constexpr uint64_t WAL_VERSION_NUMBER = 2;
Expand Down Expand Up @@ -505,13 +507,38 @@ void WriteAheadLog::WriteRowGroupData(const PersistentCollectionData &data) {
serializer.WriteProperty(101, "row_group_data", data);
serializer.End();

// only the checkpoint WAL has a checkpoint iteration - the running checkpoint does not include this commit
if (checkpoint_iteration.IsValid()) {
for (auto &block_id : data.GetBlockIds()) {
pending_checkpoint_blocks.push_back(block_id);
}
return;
}
// mark written blocks as checkpointed
auto &block_manager = GetDatabase().GetStorageManager().GetBlockManager();
for (auto &block_id : data.GetBlockIds()) {
block_manager.MarkBlockAsCheckpointed(block_id);
}
}

void WriteAheadLog::RemovePendingCheckpointBlocks(const unordered_set<block_id_t> &block_ids) {
if (block_ids.empty() || pending_checkpoint_blocks.empty()) {
return;
}
pending_checkpoint_blocks.erase(
std::remove_if(pending_checkpoint_blocks.begin(), pending_checkpoint_blocks.end(),
[&](block_id_t block_id) { return block_ids.find(block_id) != block_ids.end(); }),
pending_checkpoint_blocks.end());
}

void WriteAheadLog::MarkPendingBlocksAsCheckpointed() {
auto &block_manager = GetDatabase().GetStorageManager().GetBlockManager();
for (auto &block_id : pending_checkpoint_blocks) {
block_manager.MarkBlockAsCheckpointed(block_id);
}
pending_checkpoint_blocks.clear();
}

void WriteAheadLog::WriteDelete(DataChunk &chunk) {
D_ASSERT(chunk.size() > 0);
D_ASSERT(chunk.ColumnCount() == 1 && chunk.data[0].GetType() == LogicalType::ROW_TYPE);
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,93 @@
# name: test/sql/storage/optimistic_write/optimistic_write_concurrent_checkpoint_block_leak.test_slow
# description: Optimistically written blocks committed to the checkpoint WAL must stay free in that checkpoint's header, otherwise WAL replay leaks them.
# group: [optimistic_write]

statement ok
PRAGMA disable_checkpoint_on_shutdown

statement ok
SET checkpoint_threshold = '1TB'

statement ok
SET checkpoint_on_detach = 'disabled'

# Force row groups to flush to disk immediately, so that the WAL stores block pointers (WriteRowGroupData).
statement ok
SET write_buffer_row_group_count = 1

statement ok
ATTACH '{TEST_DIR}/concurrent_checkpoint_block_leak.db' AS db (BLOCK_SIZE 16384, ROW_GROUP_SIZE 20480)

statement ok
CREATE TABLE db.t (i BIGINT, s VARCHAR)

statement ok
INSERT INTO db.t SELECT i, repeat(md5(i::VARCHAR), 4) FROM range(50000) tt(i)

statement ok
CHECKPOINT db

# An open write transaction prevents the vacuum lock, so the checkpoint below is a concurrent checkpoint.
statement ok con_holder
BEGIN TRANSACTION

statement ok con_holder
INSERT INTO db.t VALUES (-2, 'holder')

# Make the main WAL non-empty so the checkpoint creates a checkpoint WAL.
statement ok
INSERT INTO db.t VALUES (-1, 'trigger')

statement ok
SET debug_checkpoint_sleep_ms = 3000

# Thread 0: checkpoint sleeps after WALStartCheckpoint, before writing the header.
# Thread 1: optimistic insert that commits into the checkpoint WAL during that window.
concurrentloop threadid 0 2

onlyif threadid=0
statement ok
CHECKPOINT db

onlyif threadid=1
statement ok
SELECT sleep_ms(500)

onlyif threadid=1
statement ok
INSERT INTO db.t SELECT i, repeat(md5(i::VARCHAR), 4) FROM range(50000, 150000) tt(i)

endloop

statement ok
SET debug_checkpoint_sleep_ms = 0

statement ok con_holder
ROLLBACK

# Reattach without checkpointing: the former checkpoint WAL is replayed on top of the concurrent checkpoint.
statement ok
DETACH db

statement ok
ATTACH '{TEST_DIR}/concurrent_checkpoint_block_leak.db' AS db

query I
SELECT COUNT(*) FROM db.t
----
150001

statement ok
SET debug_verify_blocks = true

statement ok
CHECKPOINT db

statement ok
DROP TABLE db.t

statement ok
CHECKPOINT db

statement ok
DETACH db
Original file line number Diff line number Diff line change
@@ -0,0 +1,116 @@
# name: test/sql/storage/optimistic_write/optimistic_write_concurrent_checkpoint_revert_commit.test_slow
# description: A failed optimistic commit into the checkpoint WAL is reverted without leaking or double-counting blocks.
# group: [optimistic_write]

statement ok
PRAGMA disable_checkpoint_on_shutdown

statement ok
SET checkpoint_threshold = '1TB'

statement ok
SET checkpoint_on_detach = 'disabled'

# Force row groups to flush to disk immediately, so that the WAL stores block pointers (WriteRowGroupData).
statement ok
SET write_buffer_row_group_count = 1

statement ok
ATTACH '{TEST_DIR}/concurrent_checkpoint_revert_commit.db' AS db (BLOCK_SIZE 16384, ROW_GROUP_SIZE 20480)

statement ok
CREATE TABLE db.t (i BIGINT, s VARCHAR)

statement ok
INSERT INTO db.t SELECT i, repeat(md5(i::VARCHAR), 4) FROM range(50000) tt(i)

statement ok
CREATE TABLE db.other (i INTEGER)

statement ok
CHECKPOINT db

# An open transaction that inserted data holds the vacuum lock, so the checkpoint below is a concurrent checkpoint.
statement ok con_holder
BEGIN TRANSACTION

statement ok con_holder
INSERT INTO db.other VALUES (1)

# Make the main WAL non-empty so the checkpoint creates a checkpoint WAL.
statement ok
INSERT INTO db.t VALUES (-1, 'trigger')

statement ok
SET debug_checkpoint_sleep_ms = 5000

# Thread 0: checkpoint sleeps after WALStartCheckpoint, before writing the header.
# Thread 1: within that window, an optimistic insert fails after writing to the checkpoint WAL and is reverted,
# then a second optimistic insert commits into the checkpoint WAL.
concurrentloop threadid 0 2

onlyif threadid=0
statement ok
CHECKPOINT db

onlyif threadid=1
statement ok
SELECT sleep_ms(500)

onlyif threadid=1
statement ok
SET debug_force_commit_failure = true

onlyif threadid=1
statement error
INSERT INTO db.t SELECT i, repeat(md5(i::VARCHAR), 4) FROM range(1000000, 1100000) tt(i)
----
Forced commit failure

onlyif threadid=1
statement ok
SET debug_force_commit_failure = false

onlyif threadid=1
statement ok
INSERT INTO db.t SELECT i, repeat(md5(i::VARCHAR), 4) FROM range(50000, 150000) tt(i)

endloop

statement ok
SET debug_checkpoint_sleep_ms = 0

statement ok con_holder
ROLLBACK

query II
SELECT COUNT(*), SUM(i) FROM db.t
----
150001 11249924999

# Reattach without checkpointing: the former checkpoint WAL must contain only the successful insert.
statement ok
DETACH db

statement ok
ATTACH '{TEST_DIR}/concurrent_checkpoint_revert_commit.db' AS db

query II
SELECT COUNT(*), SUM(i) FROM db.t
----
150001 11249924999

statement ok
SET debug_verify_blocks = true

statement ok
CHECKPOINT db

statement ok
DROP TABLE db.t

statement ok
CHECKPOINT db

statement ok
DETACH db
Loading