From 2c2f27e782e95bf13d6a8effef6779c3f5c8ac6d Mon Sep 17 00:00:00 2001 From: Xin Liao Date: Fri, 7 Aug 2026 01:39:35 +0800 Subject: [PATCH] [fix](s3) Keep the response stream usable when an error body overflows the read buffer ### What problem does this PR solve? Problem Summary: Reading from object storage fails from time to time with ``` [INTERNAL_ERROR]failed to read from : Failed to flush response stream (eof: 0, bad: 1) code=-1 type=1, request_id=failed to read ``` and succeeds when the same statement is run again. It has been hit by queries reading a rowset, by compaction, by an outfile export and by the download of an inverted index, always on an object storage that was answering `429` or `503` at that moment. `S3ObjStorageClient::get_object()` hands the buffer of the caller to the SDK as the response stream of the request, sized exactly like the requested range. The SDK writes the body of every response into that stream, the body of an error response included. The XML document of a `429 SlowDown` is a few hundred bytes, so a small ranged read cannot hold it - the read of the footer of a packed file asks for 12 bytes. `PreallocatedStreamBuf` does not implement `overflow()`, so the stream turns bad, the write callback of curl reports a short write and curl aborts the transfer with `CURLE_WRITE_ERROR`. The status code of the response is lost from there on: `CurlHttpClient::MakeRequest()` reads `CURLINFO_RESPONSE_CODE` only when curl succeeded, so the code stays at `REQUEST_NOT_MADE` (-1), and the flush check at the end of the same function replaces the retryable `NETWORK_CONNECTION` classification with `INTERNAL_FAILURE` (1). `S3CustomRetryStrategy::ShouldRetry()` declines to retry an error classified that way, and so does `S3FileReader::read_at_impl()`, which retries on `429` alone. A throttling error the server asked us to retry cancels the statement of the user instead, which is why running it again works. This also means the error carries no evidence of what really happened: the code of the response, the exception name and the request id of the object storage are all gone by the time the message is built. The fix lets the response stream grow: the body is written into the buffer of the caller as long as it fits, which is the case for every successful ranged read and keeps that path free of copies, and the remainder spills into a buffer of the stream itself, truncated at 1MB because only error documents are expected to overflow. The stream never turns bad, so curl completes the transfer, the SDK records the real status code and parses the error out of the body, and both the retry of the SDK and the retry of `S3FileReader` on `429` work again. A server or a proxy answering a ranged read with the whole object overflows the buffer as well. Such a read is still rejected, by the length check that follows the request, and now with a message that says so. Two misleading messages are fixed along the way: - `request_id=failed to read` is not a request id of the object storage. It is the string `S3FileReader` appended behind the empty request id of a failure raised by the client itself. The append is dropped and an empty request id is printed as ``. - The message of a failed read named neither the bucket nor the offset, leaving `failed to read from :` in the log whenever the key was empty. ### Release note None ### Check List (For Author) - Test: Unit Test - `be/test/io/fs/s3_response_stream_test.cpp` covers a body that fits, an error body overflowing in one write, across writes and character by character, the truncation of an oversized body, the rewind the SDK does before parsing an error, and an empty body. - Not tested end to end against a rate limited object storage. - Behavior changed: No - Does this need documentation: No --- be/src/io/fs/err_utils.cpp | 9 +- be/src/io/fs/s3_common.h | 127 +++++++++++++++++- be/src/io/fs/s3_file_reader.cpp | 7 +- be/src/io/fs/s3_obj_storage_client.cpp | 13 +- be/test/io/fs/s3_response_stream_test.cpp | 151 ++++++++++++++++++++++ 5 files changed, 296 insertions(+), 11 deletions(-) create mode 100644 be/test/io/fs/s3_response_stream_test.cpp diff --git a/be/src/io/fs/err_utils.cpp b/be/src/io/fs/err_utils.cpp index 96ac7b817e9d56..4a820e3c39c052 100644 --- a/be/src/io/fs/err_utils.cpp +++ b/be/src/io/fs/err_utils.cpp @@ -124,19 +124,22 @@ Status localfs_error(int posix_errno, std::string_view msg) { Status s3fs_error(const Aws::S3::S3Error& err, std::string_view msg) { using namespace Aws::Http; + // A failure raised by the client itself carries no request id. Printing nothing leaves a + // dangling `request_id=` that has been read as a request id of the object storage. + std::string request_id = err.GetRequestId().empty() ? "" : err.GetRequestId().c_str(); switch (err.GetResponseCode()) { case HttpResponseCode::NOT_FOUND: return Status::Error("{}: {} {} code=NOT_FOUND, type={}, request_id={}", msg, err.GetExceptionName(), err.GetMessage(), - err.GetErrorType(), err.GetRequestId()); + err.GetErrorType(), request_id); case HttpResponseCode::FORBIDDEN: return Status::Error( "{}: {} {} code=FORBIDDEN, type={}, request_id={}", msg, err.GetExceptionName(), - err.GetMessage(), err.GetErrorType(), err.GetRequestId()); + err.GetMessage(), err.GetErrorType(), request_id); default: return Status::Error( "{}: {} {} code={} type={}, request_id={}", msg, err.GetExceptionName(), - err.GetMessage(), err.GetResponseCode(), err.GetErrorType(), err.GetRequestId()); + err.GetMessage(), err.GetResponseCode(), err.GetErrorType(), request_id); } } diff --git a/be/src/io/fs/s3_common.h b/be/src/io/fs/s3_common.h index 2f42045522783c..2907e6c3fe2312 100644 --- a/be/src/io/fs/s3_common.h +++ b/be/src/io/fs/s3_common.h @@ -20,6 +20,11 @@ #include #include +#include +#include +#include +#include + namespace doris { // A non-copying iostream. @@ -34,12 +39,132 @@ class StringViewStream : Aws::Utils::Stream::PreallocatedStreamBuf, public std:: std::iostream(this) {} }; +// The AWS SDK writes the body of every response into the stream built by the response +// stream factory of the request, whatever the status of that response is. Reading an +// object range straight into the buffer of the caller therefore breaks as soon as the +// server answers with an error: the XML body of a `429 SlowDown` is a few hundred bytes +// and does not fit into the buffer of a small range read. `PreallocatedStreamBuf` does not +// implement `overflow()`, so the stream turns bad, curl aborts the transfer with +// `CURLE_WRITE_ERROR`, and the SDK reports an `INTERNAL_FAILURE` named "Failed to flush +// response stream" while never recording the status code of the response. Both the retry +// strategy of the SDK and the retry of `S3FileReader` key on that status code, so an error +// the server asked us to retry ends up cancelling the query instead. +// +// This stream buffer writes into the buffer of the caller as long as the body fits, which +// is the case for every successful ranged read, and spills the rest into a buffer of its +// own. The stream never turns bad, so the SDK reports the real status code and can parse +// the error out of the body. +class ResponseStreamBuf final : public std::streambuf { +public: + // Bodies beyond this size are truncated. Only error documents are expected to overflow + // and their leading bytes already carry the error code and the message. + static constexpr size_t MAX_SPILL_SIZE = 1024 * 1024; + + ResponseStreamBuf(void* buf, size_t nbytes) : _buf(static_cast(buf)) { + setp(_buf, _buf + nbytes); + setg(_buf, _buf, _buf); + } + +protected: + std::streamsize xsputn(const char* s, std::streamsize n) override { + if (!_spilled) { + if (n <= epptr() - pptr()) { + std::memcpy(pptr(), s, n); + pbump(static_cast(n)); + return n; + } + _spill_over(); + } + auto writable = std::min(static_cast(n), MAX_SPILL_SIZE - _spill.size()); + _spill.insert(_spill.end(), s, s + writable); + // Always report the whole write as consumed. A short write is what makes curl + // abort the transfer and lose the status code of the response. + return n; + } + + int_type overflow(int_type ch) override { + if (traits_type::eq_int_type(ch, traits_type::eof())) { + return traits_type::not_eof(ch); + } + auto c = traits_type::to_char_type(ch); + xsputn(&c, 1); + return ch; + } + + int_type underflow() override { + _reset_get_area(_read_pos()); + if (gptr() == egptr()) { + return traits_type::eof(); + } + return traits_type::to_int_type(*gptr()); + } + + pos_type seekoff(off_type off, std::ios_base::seekdir dir, + std::ios_base::openmode which) override { + auto size = static_cast(_written()); + if (which & std::ios_base::out) { + // The SDK only asks for the write position, to tell an empty body apart from a + // body it has to parse. Moving the write pointer is not supported. + return dir == std::ios_base::cur && off == 0 ? pos_type(size) : pos_type(off_type(-1)); + } + off_type pos = off; + if (dir == std::ios_base::cur) { + pos += static_cast(_read_pos()); + } else if (dir == std::ios_base::end) { + pos += size; + } + if (pos < 0 || pos > size) { + return pos_type(off_type(-1)); + } + _reset_get_area(static_cast(pos)); + return pos_type(pos); + } + + pos_type seekpos(pos_type pos, std::ios_base::openmode which) override { + return seekoff(pos, std::ios_base::beg, which); + } + +private: + // Moves what has been written so far into the spill buffer, so that the body stays + // contiguous and the SDK can parse the error out of it. + void _spill_over() { + _spill.assign(_buf, pptr()); + setp(nullptr, nullptr); + _spilled = true; + } + + // Bytes of the body held by this buffer, truncation excluded. + size_t _written() const { return _spilled ? _spill.size() : pptr() - _buf; } + + // Both areas start at the same logical offset, so the read position survives a spill. + size_t _read_pos() const { return gptr() - eback(); } + + void _reset_get_area(size_t pos) { + char* begin = _spilled ? _spill.data() : _buf; + auto size = _written(); + pos = std::min(pos, size); + setg(begin, begin + pos, begin + size); + } + + char* _buf; + std::vector _spill; + bool _spilled = false; +}; + +class ResponseStream final : public std::iostream { +public: + ResponseStream(void* buf, size_t nbytes) : std::iostream(&_buf), _buf(buf, nbytes) {} + +private: + ResponseStreamBuf _buf; +}; + // By default, the AWS SDK reads object data into an auto-growing StringStream. // To avoid copies, read directly into our preallocated buffer instead. // See https://github.com/aws/aws-sdk-cpp/issues/64 for an alternative but // functionally similar recipe. inline Aws::IOStreamFactory AwsWriteableStreamFactory(void* buf, int64_t nbytes) { - return [=]() { return Aws::New("", buf, nbytes); }; + return [=]() { return Aws::New("", buf, static_cast(nbytes)); }; } } // namespace doris diff --git a/be/src/io/fs/s3_file_reader.cpp b/be/src/io/fs/s3_file_reader.cpp index 8a6e5c0fdc4978..e4d26fab3ad908 100644 --- a/be/src/io/fs/s3_file_reader.cpp +++ b/be/src/io/fs/s3_file_reader.cpp @@ -183,9 +183,10 @@ Status S3FileReader::read_at_impl(size_t offset, Slice result, size_t* bytes_rea total_sleep_time += wait_time; continue; } else { - // Handle other errors - return std::move(Status(resp.status.code, std::move(resp.status.msg)) - .append("failed to read")); + // Handle other errors. The message already tells what failed to be read, + // appending to it only leaves a trailing token behind the request id of the + // object storage, which has been read as a request id more than once. + return {resp.status.code, std::move(resp.status.msg)}; } } if (*bytes_read != bytes_req) { diff --git a/be/src/io/fs/s3_obj_storage_client.cpp b/be/src/io/fs/s3_obj_storage_client.cpp index 0c0b0370f8097f..61da9b5d264226 100644 --- a/be/src/io/fs/s3_obj_storage_client.cpp +++ b/be/src/io/fs/s3_obj_storage_client.cpp @@ -310,17 +310,22 @@ ObjectStorageResponse S3ObjStorageClient::get_object(const ObjectStoragePathOpti if (!outcome.IsSuccess()) { record_s3_request_failed(outcome.GetError()); return {convert_to_obj_response(s3fs_error( - outcome.GetError(), fmt::format("failed to read from {}", opts.key))), + outcome.GetError(), + fmt::format("failed to read from bucket={} key={} offset={} size={}", + opts.bucket, opts.key, offset, bytes_read))), static_cast(outcome.GetError().GetResponseCode()), outcome.GetError().GetRequestId()}; } *size_return = outcome.GetResult().GetContentLength(); - // case for incomplete read + // case for incomplete read, and the case of a server or a proxy answering a ranged read + // with the whole object, which no longer fits into the buffer of the caller SYNC_POINT_CALLBACK("s3_obj_storage_client::get_object", size_return); if (*size_return != bytes_read) { return {convert_to_obj_response(Status::InternalError( - "failed to read from {}(bytes read: {}, bytes req: {}), request_id: {}", opts.key, - *size_return, bytes_read, outcome.GetResult().GetRequestId()))}; + "failed to read from bucket={} key={} offset={}(bytes read: {}, bytes req: {}), " + "request_id: {}", + opts.bucket, opts.key, offset, *size_return, bytes_read, + outcome.GetResult().GetRequestId()))}; } return ObjectStorageResponse::OK(); } diff --git a/be/test/io/fs/s3_response_stream_test.cpp b/be/test/io/fs/s3_response_stream_test.cpp new file mode 100644 index 00000000000000..d83b140ab04adf --- /dev/null +++ b/be/test/io/fs/s3_response_stream_test.cpp @@ -0,0 +1,151 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +#include + +#include +#include +#include + +#include "io/fs/s3_common.h" + +namespace doris { + +namespace { + +// What the SDK does with the body of a response it has to build an error from. +std::string drain(std::iostream& stream) { + std::stringstream out; + out << stream.rdbuf(); + return out.str(); +} + +// The XML body a MinIO answers a throttled ranged read with, shortened. +constexpr char SLOW_DOWN_BODY[] = + R"(SlowDownPlease )" + R"(reduce your request rate.data/packed_file/2666/x.bin)"; + +} // namespace + +// A body of the requested size lands in the buffer of the caller, without a copy. +TEST(ResponseStreamTest, BodyFits) { + std::string body(64, 'a'); + std::vector buffer(body.size()); + + ResponseStream stream(buffer.data(), buffer.size()); + stream.write(body.data(), body.size()); + stream.flush(); + + EXPECT_FALSE(stream.fail()); + EXPECT_EQ(body, std::string(buffer.data(), buffer.size())); + EXPECT_EQ(static_cast(body.size()), stream.tellp()); + EXPECT_EQ(body, drain(stream)); +} + +// An error body larger than the range of the read leaves the stream usable, which is what +// keeps curl from aborting the transfer and the SDK from losing the status code. +TEST(ResponseStreamTest, ErrorBodyOverflowsInOneWrite) { + std::string body(SLOW_DOWN_BODY); + // A read of the footer of a packed file is far smaller than the error document. + std::vector buffer(12); + + ResponseStream stream(buffer.data(), buffer.size()); + stream.write(body.data(), body.size()); + stream.flush(); + + EXPECT_FALSE(stream.fail()); + EXPECT_EQ(static_cast(body.size()), stream.tellp()); + EXPECT_EQ(body, drain(stream)); +} + +// curl hands the body over in chunks, so the overflow can happen in the middle of one. +TEST(ResponseStreamTest, ErrorBodyOverflowsAcrossWrites) { + std::string body(SLOW_DOWN_BODY); + std::vector buffer(16); + + ResponseStream stream(buffer.data(), buffer.size()); + size_t chunk = 7; + for (size_t pos = 0; pos < body.size(); pos += chunk) { + stream.write(body.data() + pos, std::min(chunk, body.size() - pos)); + } + stream.flush(); + + EXPECT_FALSE(stream.fail()); + EXPECT_EQ(static_cast(body.size()), stream.tellp()); + // The bytes written before the overflow are kept, so the body stays contiguous. + EXPECT_EQ(body, drain(stream)); +} + +// A body written one character at a time goes through overflow() instead of xsputn(). +TEST(ResponseStreamTest, ErrorBodyOverflowsCharByChar) { + std::string body(SLOW_DOWN_BODY); + std::vector buffer(4); + + ResponseStream stream(buffer.data(), buffer.size()); + for (char c : body) { + stream.put(c); + } + stream.flush(); + + EXPECT_FALSE(stream.fail()); + EXPECT_EQ(body, drain(stream)); +} + +// A server answering a ranged read with the whole object must not blow up the memory of the +// backend. The body is truncated, the stream stays good and the read is rejected later on by +// the length check of the caller. +TEST(ResponseStreamTest, OversizedBodyIsTruncated) { + std::string body(ResponseStreamBuf::MAX_SPILL_SIZE + 4096, 'x'); + std::vector buffer(8); + + ResponseStream stream(buffer.data(), buffer.size()); + stream.write(body.data(), body.size()); + stream.flush(); + + EXPECT_FALSE(stream.fail()); + EXPECT_EQ(static_cast(ResponseStreamBuf::MAX_SPILL_SIZE), stream.tellp()); + EXPECT_EQ(ResponseStreamBuf::MAX_SPILL_SIZE, drain(stream).size()); +} + +// The SDK rewinds the body before parsing an error out of it. +TEST(ResponseStreamTest, SeekBackAndForth) { + std::string body(SLOW_DOWN_BODY); + std::vector buffer(12); + + ResponseStream stream(buffer.data(), buffer.size()); + stream.write(body.data(), body.size()); + + EXPECT_EQ(body, drain(stream)); + stream.clear(); + stream.seekg(0); + EXPECT_EQ(body, drain(stream)); + + stream.clear(); + stream.seekg(2); + EXPECT_EQ(body.substr(2), drain(stream)); +} + +// An empty body is what tells the SDK to build the error out of the status code alone. +TEST(ResponseStreamTest, EmptyBody) { + std::vector buffer(16); + ResponseStream stream(buffer.data(), buffer.size()); + + EXPECT_EQ(std::streampos(0), stream.tellp()); + EXPECT_TRUE(drain(stream).empty()); +} + +} // namespace doris