diff options
| author | Mikhail Kopylov <[email protected]> | 2026-07-17 10:08:47 +0300 |
|---|---|---|
| committer | GitHub <[email protected]> | 2026-07-17 10:08:47 +0300 |
| commit | cf8d173d2d06da1e83a271ebcf20a8b51ae4b3c4 (patch) | |
| tree | e7c6dd860bbb013826e942efa466e405ba1a34ef | |
| parent | 9eb95d9ac18eef1b72e4e4fbdf884da3fb2bab94 (diff) | |
retry read from S3 on partial file error YQ-5501 (#46678)
5 files changed, 132 insertions, 19 deletions
diff --git a/ydb/core/kqp/ut/federated_query/s3/kqp_federated_query_ut.cpp b/ydb/core/kqp/ut/federated_query/s3/kqp_federated_query_ut.cpp index 2a9d08fcffe..1a1bdab9bf1 100644 --- a/ydb/core/kqp/ut/federated_query/s3/kqp_federated_query_ut.cpp +++ b/ydb/core/kqp/ut/federated_query/s3/kqp_federated_query_ut.cpp @@ -5,6 +5,7 @@ #include <ydb/core/kqp/common/simple/services.h> #include <ydb/core/kqp/federated_query/kqp_federated_query_helpers.h> #include <ydb/core/kqp/ut/common/kqp_ut_common.h> +#include <ydb/core/wrappers/ut_helpers/s3_mock.h> #include <ydb/public/sdk/cpp/include/ydb-cpp-sdk/client/draft/ydb_scripting.h> #include <ydb/public/sdk/cpp/include/ydb-cpp-sdk/client/operation/operation.h> #include <ydb/public/sdk/cpp/include/ydb-cpp-sdk/client/proto/accessor.h> @@ -14,6 +15,7 @@ #include <yql/essentials/utils/log/log.h> #include <library/cpp/protobuf/interop/cast.h> +#include <library/cpp/testing/common/network.h> #include <fmt/format.h> @@ -3832,6 +3834,57 @@ Y_UNIT_TEST_SUITE(KqpFederatedQuery) { UNIT_ASSERT_VALUES_EQUAL_C(result.GetStatus(), EStatus::BAD_REQUEST, result.GetIssues().ToString()); } } + + Y_UNIT_TEST(TestSuccessfulReadAfterPartialFileError) { + using NKikimr::NWrappers::NTestHelpers::TS3Mock; + + constexpr TStringBuf bucket = "partial-read-bucket"; + constexpr TStringBuf object = "test.json"; + const TString objectPath = TStringBuilder() << "/" << bucket << "/" << object; + const auto port = NTesting::GetFreePort(); + + THashMap<TString, TString> data{{objectPath, TString(TEST_CONTENT)}}; + TS3Mock s3Mock( + std::move(data), + TS3Mock::TSettings(port).WithPartialReadFailure(objectPath)); + UNIT_ASSERT_C(s3Mock.Start(), s3Mock.GetError()); + + auto kikimr = NTestUtils::MakeKikimrRunner(); + auto db = kikimr->GetQueryClient(); + + auto result = db.ExecuteQuery(fmt::format(R"( + CREATE EXTERNAL DATA SOURCE `/Root/partial_read_source` WITH ( + SOURCE_TYPE = "ObjectStorage", + LOCATION = "http://localhost:{port}/{bucket}/", + AUTH_METHOD = "NONE" + ); + CREATE EXTERNAL TABLE `/Root/partial_read_table` ( + key Utf8 NOT NULL, + value Utf8 NOT NULL + ) WITH ( + DATA_SOURCE = "/Root/partial_read_source", + LOCATION = "{object}", + FORMAT = "json_each_row" + );)", + "port"_a = port, + "bucket"_a = bucket, + "object"_a = object + ), TTxControl::NoTx()).ExtractValueSync(); + UNIT_ASSERT_VALUES_EQUAL_C(result.GetStatus(), EStatus::SUCCESS, result.GetIssues().ToOneLineString()); + + result = db.ExecuteQuery(R"( + SELECT COUNT(*) + FROM `/Root/partial_read_table`; + )", TTxControl::NoTx()).ExtractValueSync(); + UNIT_ASSERT_VALUES_EQUAL_C(result.GetStatus(), EStatus::SUCCESS, result.GetIssues().ToOneLineString()); + UNIT_ASSERT_VALUES_EQUAL(result.GetResultSets().size(), 1); + + TResultSetParser parser(result.GetResultSet(0)); + UNIT_ASSERT(parser.TryNextRow()); + UNIT_ASSERT_VALUES_EQUAL(parser.ColumnParser(0).GetUint64(), 2); + UNIT_ASSERT_VALUES_EQUAL(s3Mock.GetPartialReadFailureCount(), 1); + UNIT_ASSERT_GE(s3Mock.GetPartialReadRequestCount(), 2); + } } } // namespace NKikimr::NKqp diff --git a/ydb/core/kqp/ut/federated_query/s3/ya.make b/ydb/core/kqp/ut/federated_query/s3/ya.make index 6e855f9ceed..58f22197b5e 100644 --- a/ydb/core/kqp/ut/federated_query/s3/ya.make +++ b/ydb/core/kqp/ut/federated_query/s3/ya.make @@ -19,8 +19,10 @@ SRCS( PEERDIR( contrib/libs/aws-sdk-cpp/aws-cpp-sdk-s3 library/cpp/protobuf/interop + library/cpp/testing/common ydb/core/kqp/ut/common ydb/core/kqp/ut/federated_query/common + ydb/core/wrappers/ut_helpers ydb/library/testlib/s3_recipe_helper ydb/library/yql/providers/s3/actors ydb/public/sdk/cpp/src/client/types/operation diff --git a/ydb/core/wrappers/ut_helpers/s3_mock.cpp b/ydb/core/wrappers/ut_helpers/s3_mock.cpp index 1ce283ae906..78c678cba43 100644 --- a/ydb/core/wrappers/ut_helpers/s3_mock.cpp +++ b/ydb/core/wrappers/ut_helpers/s3_mock.cpp @@ -17,6 +17,7 @@ namespace NTestHelpers { TS3Mock::TSettings::TSettings() : CorruptETags(false) , RejectUploadParts(false) + , PartialReadFailures(0) { } @@ -24,6 +25,7 @@ TS3Mock::TSettings::TSettings(ui16 port) : HttpOptions(THttpServer::TOptions(port).SetThreads(1)) , CorruptETags(false) , RejectUploadParts(false) + , PartialReadFailures(0) { } @@ -42,6 +44,12 @@ TS3Mock::TSettings& TS3Mock::TSettings::WithRejectUploadParts(bool value) { return *this; } +TS3Mock::TSettings& TS3Mock::TSettings::WithPartialReadFailure(TString path, ui32 count) { + PartialReadPath = std::move(path); + PartialReadFailures = count; + return *this; +} + TS3Mock::TRequest::EMethod TS3Mock::TRequest::ParseMethod(const char* str) { if (strnicmp(str, "HEAD", 4) == 0) { return EMethod::Head; @@ -129,20 +137,30 @@ bool TS3Mock::TRequest::HttpServeRead(const TReplyParams& params, EMethod method ShuffleRange(etag); } - params.Output << "HTTP/1.1 200 Ok\r\n"; + params.Output << (rangeHeader ? "HTTP/1.1 206 Partial Content\r\n" : "HTTP/1.1 200 Ok\r\n"); THttpHeaders headers; headers.AddHeader("ETag", etag); if (method == EMethod::Get) { - headers.AddHeader("Content-Length", range.second - range.first + 1); + const size_t responseSize = range.second - range.first + 1; + const bool failPartialRead = Parent->ShouldFailPartialRead(path); + headers.AddHeader("Content-Length", responseSize); headers.AddHeader("Content-Type", "application/octet-stream"); - headers.OutTo(¶ms.Output); if (rangeHeader) { headers.AddHeader("Accept-Ranges", "bytes"); headers.AddHeader("Content-Range", "bytes " + ::ToString(range.first) + "-" + ToString(range.second) + "/" + ::ToString(content.size())); } + if (failPartialRead) { + params.Output.EnableKeepAlive(false); + headers.AddHeader("Connection", "close"); + } + headers.OutTo(¶ms.Output); params.Output << "\r\n"; - params.Output << content.SubStr(range.first, range.second - range.first + 1); + params.Output << content.SubStr(range.first, failPartialRead ? responseSize / 2 : responseSize); + if (failPartialRead) { + params.Output.Finish(); + return true; + } } else { headers.AddHeader("Content-Length", content.size()); headers.OutTo(¶ms.Output); @@ -153,31 +171,33 @@ bool TS3Mock::TRequest::HttpServeRead(const TReplyParams& params, EMethod method return true; } -TString BuildContentXML(const TString& path) { +TString BuildContentXML(const TString& path, size_t size) { return Sprintf(R"( <Contents> <Key>%s</Key> + <Size>%zu</Size> </Contents> - )", path.c_str()); + )", path.c_str(), size); } -TString BuildContentListXML(const TVector<TString>& paths) { +TString BuildContentListXML(const TVector<std::pair<TString, size_t>>& objects) { TString result; - for (const auto& path : paths) { - result += BuildContentXML(path); + for (const auto& [path, size] : objects) { + result += BuildContentXML(path, size); } return result; } -TString BuildListObjectsXML(const TVector<TString>& paths, const TStringBuf bucketName) { - return Sprintf(R"( - <?xml version="1.0" encoding="UTF-8"?> - <ListBucketResult> +TString BuildListObjectsXML(const TVector<std::pair<TString, size_t>>& objects, const TStringBuf bucketName) { + return Sprintf(R"(<?xml version="1.0" encoding="UTF-8"?> + <ListBucketResult xmlns="http://s3.amazonaws.com/doc/2006-03-01/"> <Name>%s</Name> <IsTruncated>false</IsTruncated> + <MaxKeys>%zu</MaxKeys> + <KeyCount>%zu</KeyCount> %s </ListBucketResult> - )", bucketName.data(), BuildContentListXML(paths).c_str()); + )", TString(bucketName).c_str(), objects.size(), objects.size(), BuildContentListXML(objects).c_str()); } bool TS3Mock::TRequest::HttpServeList(const TReplyParams& params, TStringBuf bucketName, const TString& prefix) { @@ -185,23 +205,26 @@ bool TS3Mock::TRequest::HttpServeList(const TReplyParams& params, TStringBuf buc params.Output << "HTTP/1.1 200 Ok\r\n"; THttpHeaders headers; - TVector<TString> paths; - const TString bucketPrefix = TStringBuilder() << bucketName << "/"; + TVector<std::pair<TString, size_t>> objects; + TString bucketPrefix(bucketName); + if (!bucketPrefix.EndsWith('/')) { + bucketPrefix += '/'; + } const TString keyPrefix = TStringBuilder() << bucketPrefix << prefix; for (const auto& [key, value] : Parent->Data) { if (key.StartsWith(keyPrefix)) { - paths.push_back(key.substr(bucketPrefix.size())); + objects.emplace_back(key.substr(bucketPrefix.size()), value.size()); } } - TString xml = BuildListObjectsXML(paths, bucketName); + TString xml = BuildListObjectsXML(objects, bucketName); headers.AddHeader("Content-Type", "application/xml"); headers.AddHeader("Content-Length", xml.length()); headers.OutTo(¶ms.Output); - params.Output << xml; params.Output << "\r\n"; + params.Output << xml; params.Output.Flush(); return true; @@ -455,6 +478,7 @@ const char* TS3Mock::GetError() { TS3Mock::TS3Mock(const TSettings& settings) : Settings(settings) + , PartialReadFailuresLeft(settings.PartialReadFailures) , HttpServer(this, settings.HttpOptions) { } @@ -462,6 +486,7 @@ TS3Mock::TS3Mock(const TSettings& settings) TS3Mock::TS3Mock(THashMap<TString, TString>&& data, const TSettings& settings) : Settings(settings) , Data(std::move(data)) + , PartialReadFailuresLeft(settings.PartialReadFailures) , HttpServer(this, settings.HttpOptions) { } @@ -469,10 +494,28 @@ TS3Mock::TS3Mock(THashMap<TString, TString>&& data, const TSettings& settings) TS3Mock::TS3Mock(const THashMap<TString, TString>& data, const TSettings& settings) : Settings(settings) , Data(data) + , PartialReadFailuresLeft(settings.PartialReadFailures) , HttpServer(this, settings.HttpOptions) { } +bool TS3Mock::ShouldFailPartialRead(TStringBuf path) { + if (path != Settings.PartialReadPath) { + return false; + } + + ++PartialReadRequestCount; + + ui32 failuresLeft = PartialReadFailuresLeft.load(); + while (failuresLeft) { + if (PartialReadFailuresLeft.compare_exchange_weak(failuresLeft, failuresLeft - 1)) { + ++PartialReadFailureCount; + return true; + } + } + return false; +} + TClientRequest* TS3Mock::CreateClient() { return new TRequest(this); } diff --git a/ydb/core/wrappers/ut_helpers/s3_mock.h b/ydb/core/wrappers/ut_helpers/s3_mock.h index 6c77cba9b8f..367649dc091 100644 --- a/ydb/core/wrappers/ut_helpers/s3_mock.h +++ b/ydb/core/wrappers/ut_helpers/s3_mock.h @@ -5,6 +5,8 @@ #include <library/cpp/http/server/http.h> #include <library/cpp/cgiparam/cgiparam.h> +#include <atomic> + #include <util/generic/hash.h> #include <util/generic/string.h> #include <util/generic/vector.h> @@ -19,6 +21,8 @@ public: THttpServer::TOptions HttpOptions; bool CorruptETags; bool RejectUploadParts; + TString PartialReadPath; + ui32 PartialReadFailures; TSettings(); explicit TSettings(ui16 port); @@ -26,6 +30,7 @@ public: TSettings& WithHttpOptions(const THttpServer::TOptions& opts); TSettings& WithCorruptETags(bool value); TSettings& WithRejectUploadParts(bool value); + TSettings& WithPartialReadFailure(TString path, ui32 count = 1); }; // TSettings @@ -75,10 +80,19 @@ public: const THashMap<TString, TString>& GetData() const override { return Data; } THashMap<TString, TString>& GetData() override { return Data; } + ui32 GetPartialReadRequestCount() const { return PartialReadRequestCount.load(); } + ui32 GetPartialReadFailureCount() const { return PartialReadFailureCount.load(); } + private: + bool ShouldFailPartialRead(TStringBuf path); + const TSettings Settings; THashMap<TString, TString> Data; + std::atomic<ui32> PartialReadFailuresLeft; + std::atomic<ui32> PartialReadRequestCount = 0; + std::atomic<ui32> PartialReadFailureCount = 0; + int NextUploadId = 1; THashMap<std::pair<TString, TString>, TVector<TString>> MultipartUploads; THttpServer HttpServer; diff --git a/ydb/library/yql/providers/common/http_gateway/yql_http_default_retry_policy.cpp b/ydb/library/yql/providers/common/http_gateway/yql_http_default_retry_policy.cpp index 4b91e12a8f2..3d1d7b9722c 100644 --- a/ydb/library/yql/providers/common/http_gateway/yql_http_default_retry_policy.cpp +++ b/ydb/library/yql/providers/common/http_gateway/yql_http_default_retry_policy.cpp @@ -15,6 +15,7 @@ std::unordered_set<CURLcode> FqRetriedCurlCodes() { CURLE_BAD_DOWNLOAD_RESUME, CURLE_SEND_ERROR, CURLE_RECV_ERROR, + CURLE_PARTIAL_FILE, CURLE_NO_CONNECTION_AVAILABLE, CURLE_GOT_NOTHING, CURLE_COULDNT_RESOLVE_HOST |
