summaryrefslogtreecommitdiffstats
diff options
context:
space:
mode:
authorMikhail Kopylov <[email protected]>2026-07-17 10:08:47 +0300
committerGitHub <[email protected]>2026-07-17 10:08:47 +0300
commitcf8d173d2d06da1e83a271ebcf20a8b51ae4b3c4 (patch)
treee7c6dd860bbb013826e942efa466e405ba1a34ef
parent9eb95d9ac18eef1b72e4e4fbdf884da3fb2bab94 (diff)
retry read from S3 on partial file error YQ-5501 (#46678)
-rw-r--r--ydb/core/kqp/ut/federated_query/s3/kqp_federated_query_ut.cpp53
-rw-r--r--ydb/core/kqp/ut/federated_query/s3/ya.make2
-rw-r--r--ydb/core/wrappers/ut_helpers/s3_mock.cpp81
-rw-r--r--ydb/core/wrappers/ut_helpers/s3_mock.h14
-rw-r--r--ydb/library/yql/providers/common/http_gateway/yql_http_default_retry_policy.cpp1
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(&params.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(&params.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(&params.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(&params.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