This is an automated email from the ASF dual-hosted git repository.
mymeiyi pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/doris.git
The following commit(s) were added to refs/heads/master by this push:
new fc3235c3c95 [fix](be) Handle incomplete WAL records during replay
(#67204)
fc3235c3c95 is described below
commit fc3235c3c950e7c8b1103c3bf513f2965c13c51d
Author: meiyi <[email protected]>
AuthorDate: Wed Sep 2 11:14:44 2026 +0800
[fix](be) Handle incomplete WAL records during replay (#67204)
Problem Summary: when be restart, wal file may be incomplete, replay wal
should skip the incomplete block data
---
be/src/load/group_commit/wal/wal_file_reader.cpp | 70 ++++++++++++++++---
be/src/load/group_commit/wal/wal_table.cpp | 3 +-
be/src/load/group_commit/wal/wal_writer.cpp | 13 ++++
be/test/format/wal/wal_reader_writer_test.cpp | 80 ++++++++++++++++++++++
.../test_group_commit_replay_wal.groovy | 3 +-
5 files changed, 157 insertions(+), 12 deletions(-)
diff --git a/be/src/load/group_commit/wal/wal_file_reader.cpp
b/be/src/load/group_commit/wal/wal_file_reader.cpp
index cd9430a6c18..c2ebede07f9 100644
--- a/be/src/load/group_commit/wal/wal_file_reader.cpp
+++ b/be/src/load/group_commit/wal/wal_file_reader.cpp
@@ -83,6 +83,13 @@ Status WalFileReader::read_block(PBlock& block) {
if (_offset >= file_reader->size()) {
return Status::EndOfFile("end of wal file");
}
+ const size_t file_size = file_reader->size();
+ if (file_size - _offset < WalWriter::LENGTH_SIZE) {
+ LOG(WARNING) << "ignore incomplete wal tail, path=" << _file_name
+ << ", file_size=" << file_size << ", read_offset=" <<
_offset
+ << ", expected_length_bytes=" << WalWriter::LENGTH_SIZE;
+ return Status::EndOfFile("end of wal file");
+ }
size_t bytes_read = 0;
uint8_t row_len_buf[WalWriter::LENGTH_SIZE];
RETURN_IF_ERROR(
@@ -92,18 +99,34 @@ Status WalFileReader::read_block(PBlock& block) {
if (block_len == 0) {
return Status::DataQualityError("fail to read wal {} ,block is empty",
_file_name);
}
- if (_offset == file_reader->size()) {
- LOG(WARNING) << "need read block with length=" << block_len << ", but
offset=" << _offset
- << " reached end of WAL (path=" << _file_name
- << ", size=" << file_reader->size() << ")";
+ const size_t remaining_bytes = file_size - _offset;
+ if (block_len > remaining_bytes) {
+ LOG(WARNING) << "ignore incomplete wal tail, path=" << _file_name
+ << ", file_size=" << file_size << ", read_offset=" <<
_offset
+ << ", block_bytes=" << block_len << ", available_bytes="
<< remaining_bytes;
return Status::EndOfFile("end of wal file");
}
// read block
std::string block_buf;
block_buf.resize(block_len);
RETURN_IF_ERROR(file_reader->read_at(_offset, {block_buf.c_str(),
block_len}, &bytes_read));
+ if (bytes_read != block_len) {
+ LOG(WARNING) << "ignore incomplete wal tail, path=" << _file_name
+ << ", file_size=" << file_size << ", read_offset=" <<
_offset
+ << ", block_bytes=" << block_len << ", read_block_bytes="
<< bytes_read;
+ return Status::EndOfFile("end of wal file");
+ }
RETURN_IF_ERROR(_deserialize(block, block_buf, block_len, bytes_read));
_offset += block_len;
+ const size_t checksum_bytes = file_size - _offset;
+ if (checksum_bytes < WalWriter::CHECKSUM_SIZE) {
+ LOG(WARNING) << "replay wal block without complete checksum, path=" <<
_file_name
+ << ", file_size=" << file_size << ", read_offset=" <<
_offset
+ << ", block_bytes=" << block_len << ", checksum_bytes="
<< checksum_bytes
+ << ", expected_checksum_bytes=" <<
WalWriter::CHECKSUM_SIZE;
+ _offset = file_size;
+ return Status::OK();
+ }
// checksum
uint8_t checksum_len_buf[WalWriter::CHECKSUM_SIZE];
RETURN_IF_ERROR(file_reader->read_at(_offset, {checksum_len_buf,
WalWriter::CHECKSUM_SIZE},
@@ -115,33 +138,60 @@ Status WalFileReader::read_block(PBlock& block) {
}
Status WalFileReader::read_header(uint32_t& version, std::string& col_ids) {
- if (file_reader->size() == 0) {
- return Status::DataQualityError("empty file");
- }
+ const size_t file_size = file_reader->size();
+ auto incomplete_header = [&](const char* field, size_t expected_bytes,
size_t actual_bytes) {
+ LOG(WARNING) << "ignore incomplete wal header, path=" << _file_name
+ << ", file_size=" << file_size << ", read_offset=" <<
_offset
+ << ", field=" << field << ", expected_bytes=" <<
expected_bytes
+ << ", actual_bytes=" << actual_bytes;
+ return Status::DataQualityError(
+ "incomplete wal header {}, field={}, expected_bytes={},
actual_bytes={}",
+ _file_name, field, expected_bytes, actual_bytes);
+ };
size_t bytes_read = 0;
+ if (file_size - _offset < k_wal_magic_length) {
+ return incomplete_header("magic", k_wal_magic_length, file_size -
_offset);
+ }
std::string magic_str;
magic_str.resize(k_wal_magic_length);
RETURN_IF_ERROR(file_reader->read_at(_offset, magic_str, &bytes_read));
+ if (bytes_read != k_wal_magic_length) {
+ return incomplete_header("magic", k_wal_magic_length, bytes_read);
+ }
if (strcmp(magic_str.c_str(), k_wal_magic) != 0) {
return Status::Corruption("Bad wal file {}: magic number not match",
_file_name);
}
_offset += k_wal_magic_length;
+ if (file_size - _offset < WalWriter::VERSION_SIZE) {
+ return incomplete_header("version", WalWriter::VERSION_SIZE, file_size
- _offset);
+ }
uint8_t version_buf[WalWriter::VERSION_SIZE];
RETURN_IF_ERROR(
file_reader->read_at(_offset, {version_buf,
WalWriter::VERSION_SIZE}, &bytes_read));
+ if (bytes_read != WalWriter::VERSION_SIZE) {
+ return incomplete_header("version", WalWriter::VERSION_SIZE,
bytes_read);
+ }
_offset += WalWriter::VERSION_SIZE;
version = decode_fixed32_le(version_buf);
+ if (file_size - _offset < WalWriter::LENGTH_SIZE) {
+ return incomplete_header("column_ids_length", WalWriter::LENGTH_SIZE,
file_size - _offset);
+ }
uint8_t len_buf[WalWriter::LENGTH_SIZE];
RETURN_IF_ERROR(file_reader->read_at(_offset, {len_buf,
WalWriter::LENGTH_SIZE}, &bytes_read));
+ if (bytes_read != WalWriter::LENGTH_SIZE) {
+ return incomplete_header("column_ids_length", WalWriter::LENGTH_SIZE,
bytes_read);
+ }
_offset += WalWriter::LENGTH_SIZE;
size_t len = decode_fixed64_le(len_buf);
+ if (len > file_size - _offset) {
+ return incomplete_header("column_ids", len, file_size - _offset);
+ }
col_ids.resize(len);
RETURN_IF_ERROR(file_reader->read_at(_offset, col_ids, &bytes_read));
- _offset += len;
if (len != bytes_read) {
- return Status::InternalError("failed to read header expected= " +
std::to_string(len) +
- ",actually=" +
std::to_string(bytes_read));
+ return incomplete_header("column_ids", len, bytes_read);
}
+ _offset += len;
return Status::OK();
}
diff --git a/be/src/load/group_commit/wal/wal_table.cpp
b/be/src/load/group_commit/wal/wal_table.cpp
index 34109affe9b..69e5d7cf2b7 100644
--- a/be/src/load/group_commit/wal/wal_table.cpp
+++ b/be/src/load/group_commit/wal/wal_table.cpp
@@ -92,7 +92,8 @@ Status WalTable::_relay_wal_one_by_one() {
int64_t file_size = 0;
std::filesystem::path file_path(wal_info->get_wal_path());
if (!std::filesystem::exists(file_path)) {
- st = Status::InternalError("wal file {} does not exist",
wal_info->get_wal_path());
+ LOG(WARNING) << "skip replay missing wal=" <<
wal_info->get_wal_path();
+ st = Status::OK();
} else {
file_size = std::filesystem::file_size(file_path);
st = _replay_wal_internal(wal_info->get_wal_path());
diff --git a/be/src/load/group_commit/wal/wal_writer.cpp
b/be/src/load/group_commit/wal/wal_writer.cpp
index 57f5e6c59c2..1989c3e588e 100644
--- a/be/src/load/group_commit/wal/wal_writer.cpp
+++ b/be/src/load/group_commit/wal/wal_writer.cpp
@@ -31,6 +31,7 @@
#include "load/group_commit/wal/wal_manager.h"
#include "runtime/cluster_info.h"
#include "storage/storage_engine.h"
+#include "util/debug_points.h"
#include "util/thrift_rpc_helper.h"
namespace doris {
@@ -135,6 +136,18 @@ Status WalWriter::append_blocks(const PBlockArray& blocks)
{
"failed to write block to wal expected= " +
std::to_string(total_size) +
",actually=" + std::to_string(offset));
}
+ DBUG_EXECUTE_IF("WalWriter.append_blocks.write_incomplete_tail", {
+ if (!blocks.empty()) {
+ uint8_t len_buf[sizeof(uint64_t)];
+ uint64_t block_length = blocks.back()->ByteSizeLong();
+ encode_fixed64_le(len_buf, block_length);
+ RETURN_IF_ERROR(_file_writer->append({len_buf, sizeof(uint64_t)}));
+
+ std::string content = blocks.back()->SerializeAsString();
+ content.resize(content.size() / 2);
+ RETURN_IF_ERROR(_file_writer->append(content));
+ }
+ });
return Status::OK();
}
diff --git a/be/test/format/wal/wal_reader_writer_test.cpp
b/be/test/format/wal/wal_reader_writer_test.cpp
index 2b6a554926d..1b5ebffe10e 100644
--- a/be/test/format/wal/wal_reader_writer_test.cpp
+++ b/be/test/format/wal/wal_reader_writer_test.cpp
@@ -17,6 +17,7 @@
#include <gen_cpp/internal_service.pb.h>
#include <gtest/gtest.h>
+#include <array>
#include <filesystem>
#include <memory>
@@ -141,4 +142,83 @@ TEST_F(WalReaderWriterTest, TestWriteAndRead1) {
static_cast<void>(wal_reader.finalize());
EXPECT_EQ(3, block_count);
}
+
+TEST_F(WalReaderWriterTest, TestReadIncompleteLastRecord) {
+ PBlock first_block;
+ PBlock last_block;
+ generate_block(first_block, 0);
+ generate_block(last_block, block_rows);
+
+ const size_t first_record_size =
+ WalWriter::LENGTH_SIZE + first_block.ByteSizeLong() +
WalWriter::CHECKSUM_SIZE;
+ const size_t last_block_end =
+ first_record_size + WalWriter::LENGTH_SIZE +
last_block.ByteSizeLong();
+ const std::array<size_t, 6> truncated_sizes = {
+ first_record_size + WalWriter::LENGTH_SIZE / 2,
+ first_record_size + WalWriter::LENGTH_SIZE +
last_block.ByteSizeLong() / 2,
+ last_block_end,
+ last_block_end + 1,
+ last_block_end + 2,
+ last_block_end + 3};
+
+ for (size_t i = 0; i < truncated_sizes.size(); ++i) {
+ std::string file_name = _s_test_data_path + "/incomplete_last_record_"
+ std::to_string(i);
+ auto wal_writer = WalWriter(file_name);
+ ASSERT_TRUE(wal_writer.init(io::global_local_filesystem()).ok());
+ ASSERT_TRUE(wal_writer.append_blocks({&first_block,
&last_block}).ok());
+ ASSERT_TRUE(wal_writer.finalize().ok());
+ ASSERT_NO_THROW(std::filesystem::resize_file(file_name,
truncated_sizes[i]));
+
+ auto wal_reader = WalFileReader(file_name);
+ ASSERT_TRUE(wal_reader.init().ok());
+ PBlock block;
+ EXPECT_TRUE(wal_reader.read_block(block).ok());
+ auto st = wal_reader.read_block(block);
+ if (i < 2) {
+ EXPECT_TRUE(st.is<ErrorCode::END_OF_FILE>());
+ } else {
+ EXPECT_TRUE(st.ok());
+ Block deserialized_block;
+ size_t uncompressed_size = 0;
+ int64_t uncompressed_time = 0;
+ EXPECT_TRUE(
+ deserialized_block.deserialize(block, &uncompressed_size,
&uncompressed_time)
+ .ok());
+ EXPECT_EQ(block_rows, deserialized_block.rows());
+
EXPECT_TRUE(wal_reader.read_block(block).is<ErrorCode::END_OF_FILE>());
+ }
+ EXPECT_TRUE(wal_reader.finalize().ok());
+ }
+}
+
+TEST_F(WalReaderWriterTest, TestReadIncompleteHeader) {
+ const std::string column_ids = "1,2";
+ const size_t version_header_size = k_wal_magic_length +
WalWriter::VERSION_SIZE;
+ const size_t fixed_header_size = version_header_size +
WalWriter::LENGTH_SIZE;
+ const std::array<size_t, 8> truncated_sizes = {0,
+ k_wal_magic_length - 1,
+ k_wal_magic_length,
+ version_header_size - 1,
+ version_header_size,
+ fixed_header_size - 1,
+ fixed_header_size,
+ fixed_header_size +
column_ids.size() - 1};
+
+ for (size_t i = 0; i < truncated_sizes.size(); ++i) {
+ std::string file_name = _s_test_data_path + "/incomplete_header_" +
std::to_string(i);
+ auto wal_writer = WalWriter(file_name);
+ ASSERT_TRUE(wal_writer.init(io::global_local_filesystem()).ok());
+ ASSERT_TRUE(wal_writer.append_header(column_ids).ok());
+ ASSERT_TRUE(wal_writer.finalize().ok());
+ ASSERT_NO_THROW(std::filesystem::resize_file(file_name,
truncated_sizes[i]));
+
+ auto wal_reader = WalFileReader(file_name);
+ ASSERT_TRUE(wal_reader.init().ok());
+ uint32_t version = 0;
+ std::string actual_column_ids;
+ EXPECT_TRUE(wal_reader.read_header(version, actual_column_ids)
+ .is<ErrorCode::DATA_QUALITY_ERROR>());
+ EXPECT_TRUE(wal_reader.finalize().ok());
+ }
+}
} // namespace doris
diff --git
a/regression-test/suites/insert_p0/group_commit/test_group_commit_replay_wal.groovy
b/regression-test/suites/insert_p0/group_commit/test_group_commit_replay_wal.groovy
index 856515b3350..3e43f333000 100644
---
a/regression-test/suites/insert_p0/group_commit/test_group_commit_replay_wal.groovy
+++
b/regression-test/suites/insert_p0/group_commit/test_group_commit_replay_wal.groovy
@@ -76,6 +76,7 @@ suite("test_group_commit_replay_wal", "nonConcurrent") {
// load fail and abort fail, wal should not be deleted and retry
try {
+
GetDebugPoint().enableDebugPointForAllBEs("WalWriter.append_blocks.write_incomplete_tail",
[execute: 1])
GetDebugPoint().enableDebugPointForAllBEs("LoadBlockQueue._finish_group_commit_load.load_error")
GetDebugPoint().enableDebugPointForAllFEs("FrontendServiceImpl.loadTxnRollback.error")
streamLoad {
@@ -172,4 +173,4 @@ suite("test_group_commit_replay_wal", "nonConcurrent") {
sleep(2000)
}
assertEquals(0, wal_count)
-}
\ No newline at end of file
+}
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]