Skip to content
Merged
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 CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,14 @@ and this project adheres to [Semantic Versioning](http://semver.org/spec/v2.0.0.

## [Unreleased]

### Fixed
Comment thread
Serpentian marked this conversation as resolved.
* Fix consistency violation during rebalancing: a request could be retried on a
storage that did not own the bucket yet, so a read returned no data and a
write went to the wrong replicaset (#509).
* Stop retrying requests blindly: a request is retried only if the router has
recovered from the error, and a retry never runs past the original timeout —
before, a request could take up to twice as long as requested.

## [1.7.5] - 15-07-26

### Added
Expand Down
197 changes: 143 additions & 54 deletions crud/common/call.lua
Original file line number Diff line number Diff line change
Expand Up @@ -53,27 +53,16 @@ function call.get_vshard_call_name(mode, prefer_replica, balance)
return 'callbre'
end

local function wrap_vshard_err(vshard_router, err, func_name, replicaset_id, bucket_id)
local function wrap_vshard_err(err, func_name, replicaset_id)
-- Do not rewrite ShardingHashMismatchError class.
if err.class_name == sharding_utils.ShardingHashMismatchError.name then
return errors.wrap(err)
end

if replicaset_id == nil then
local replicaset, _ = vshard_router:route(bucket_id)
if replicaset == nil then
return CallError:new(
"Function returned an error, but we couldn't figure out the replicaset: %s", err
)
end

replicaset_id = utils.get_replicaset_id(vshard_router, replicaset)

if replicaset_id == nil then
Comment thread
Serpentian marked this conversation as resolved.
return CallError:new(
"Function returned an error, but we couldn't figure out the replicaset id: %s", err
)
end
return CallError:new(
"Function returned an error, but we couldn't figure out the replicaset id: %s", err
)
end

err = utils.update_storage_call_error_description(err, func_name, replicaset_id)
Expand All @@ -84,47 +73,145 @@ local function wrap_vshard_err(vshard_router, err, func_name, replicaset_id, buc
))
end
Comment thread
Serpentian marked this conversation as resolved.

--- Executes a vshard call and retries once after performing recovery actions
--- like bucket cache reset, destination redirect (for single calls), or master discovery.
local function call_with_retry_and_recovery(vshard_router,
replicaset, method, func_name, func_args, call_opts, is_single_call)
--- Executes a CRUD function on a vshard replicaset.
local function call_on_replicaset(replicaset, method, func_name, func_args, call_opts)
local func_args_ext = utils.append_array({ box.session.effective_user(), func_name }, func_args)
return replicaset[method](replicaset, CRUD_CALL_FUNC_NAME, func_args_ext, call_opts)
end

-- In case cluster was just bootstrapped with auto master discovery,
-- replicaset may miss master.
local resp, err = replicaset[method](replicaset, CRUD_CALL_FUNC_NAME, func_args_ext, call_opts)

if err == nil then
return resp, err
end
--- The bucket is not on the replicaset it is routed to anymore, so its
--- route is stale and must be dropped from the router cache.
local BUCKET_MOVED_ERRS = {
WRONG_BUCKET = true,
BUCKET_IS_LOCKED = true,
TRANSFER_IS_IN_PROGRESS = true,
}

--- Performs a recovery action for a call error:
--- * MISSING_MASTER - the replicaset master is not discovered yet, so it
--- is discovered explicitly;
--- * NON_MASTER - the cached master is stale, but the bucket did not move,
--- so the master is updated and the same replicaset can be retried;
--- * WRONG_BUCKET, BUCKET_IS_LOCKED, TRANSFER_IS_IN_PROGRESS - the bucket
--- route is stale, so the route cache is reset and the bucket is routed
--- anew.
---
--- The routes of all the moved buckets are reset even when the request itself
--- cannot be retried: the cache is stale regardless of the retry decision.
---
--- Returns:
--- * replicaset - the replicaset to retry the request on;
--- * nil - the request cannot be retried;
--- * nil, err - the recovery itself failed, so the request cannot be retried.
local function recover_from_err(vshard_router, replicaset, err)
local vshard_err = err

-- This is a partial copy of error handling from vshard.router.router_call_impl()
-- It is much simpler mostly because bucket_set() can't be accessed from outside vshard.
if err.class_name == bucket_ref_unref.BucketRefError.name then
if is_single_call and #err.bucket_ref_errs == 1 then
local single_err = err.bucket_ref_errs[1]
local destination = single_err.vshard_err.destination
if destination and vshard_router.replicasets[destination] then
replicaset = vshard_router.replicasets[destination]
for _, bucket_ref_err in ipairs(err.bucket_ref_errs) do
if BUCKET_MOVED_ERRS[bucket_ref_err.vshard_err.name] then
vshard_router:_bucket_reset(bucket_ref_err.bucket_id)
end
end

for _, bucket_ref_err in pairs(err.bucket_ref_errs) do
local bucket_id = bucket_ref_err.bucket_id
local vshard_err = bucket_ref_err.vshard_err
if vshard_err.name == 'WRONG_BUCKET' or
vshard_err.name == 'BUCKET_IS_LOCKED' or
vshard_err.name == 'TRANSFER_IS_IN_PROGRESS' then
vshard_router:_bucket_reset(bucket_id)
-- A request that failed on several buckets has no single replicaset
-- to be retried on, so only a single-bucket one is recovered further.
if #err.bucket_ref_errs ~= 1 then
return nil
end

vshard_err = err.bucket_ref_errs[1].vshard_err

if BUCKET_MOVED_ERRS[vshard_err.name] then
-- The stale route has been reset above, so the bucket is routed
-- anew.
local new_replicaset, route_err = vshard_router:route(err.bucket_ref_errs[1].bucket_id)
if route_err ~= nil then
return nil, CallError:new(
"Failed to get router replicaset: %s, after an error: %s",
tostring(route_err),
tostring(err)
)
end
return new_replicaset
end
elseif err.name == 'MISSING_MASTER' and replicaset.locate_master ~= nil then
end

if vshard_err.name == 'MISSING_MASTER' then
replicaset:locate_master()
-- The master is not found, so the retry would get the same error.
if replicaset.master == nil then
return nil
end
return replicaset
elseif vshard_err.name == 'NON_MASTER' then
-- If the master update failed, the retry would get the same error.
if not replicaset:update_master(vshard_err.replica, vshard_err.master) then
return nil
end
return replicaset
end

-- Retry only once: should be enough for initial discovery,
-- otherwise force user fix up cluster bootstrap.
return replicaset[method](replicaset, CRUD_CALL_FUNC_NAME, func_args_ext, call_opts)
return nil
end

local function call_single_with_recovery(vshard_router,
replicaset, method, func_name, func_args, call_opts)
local deadline = fiber_clock() + call_opts.timeout

local resp, err = call_on_replicaset(replicaset, method, func_name, func_args, call_opts)
if err == nil then
return resp, err, replicaset.id
end

local retry_replicaset, recover_err = recover_from_err(vshard_router, replicaset, err)
if retry_replicaset == nil then
return resp, recover_err or err, replicaset.id
end

local timeout = deadline - fiber_clock()
if timeout <= 0 then
return resp, err, replicaset.id
end

replicaset = retry_replicaset

call_opts.timeout = timeout
if call_opts.request_timeout ~= nil and call_opts.request_timeout > timeout then
call_opts.request_timeout = timeout
end

resp, err = call_on_replicaset(replicaset, method, func_name, func_args, call_opts)
return resp, err, replicaset.id
end

local function call_map_with_recovery(replicaset, method, func_name, func_args,
Comment thread
Serpentian marked this conversation as resolved.
call_opts, deadline)
local future, err = call_on_replicaset(replicaset, method, func_name, func_args, call_opts)
if err == nil or err.name ~= 'MISSING_MASTER' then
return future, err
end

if fiber_clock() >= deadline then
return future, err
end

replicaset:locate_master()

-- The master is not found, so the retry would get the same error.
if replicaset.master == nil then
return future, err
end

local timeout = deadline - fiber_clock()
if timeout <= 0 then
return future, err
end

if call_opts.request_timeout ~= nil and call_opts.request_timeout > timeout then
call_opts.request_timeout = timeout
end

return call_on_replicaset(replicaset, method, func_name, func_args, call_opts)
end

function call.map(vshard_router, func_name, func_args, opts)
Expand Down Expand Up @@ -169,11 +256,13 @@ function call.map(vshard_router, func_name, func_args, opts)
is_async = true,
request_timeout = opts.mode == 'read' and opts.request_timeout or nil,
}

local deadline = fiber_clock() + timeout
while iter:has_next() do
local args, replicaset, replicaset_id = iter:get()

local future, err = call_with_retry_and_recovery(vshard_router, replicaset, vshard_call_name,
func_name, args, call_opts, false)
local future, err = call_map_with_recovery(replicaset, vshard_call_name,
func_name, args, call_opts, deadline)

if err ~= nil then
local result_info = {
Expand All @@ -195,7 +284,6 @@ function call.map(vshard_router, func_name, func_args, opts)
futures_by_replicasets[replicaset_id] = future
end

local deadline = fiber_clock() + timeout
for replicaset_id, future in pairs(futures_by_replicasets) do
local wait_timeout = deadline - fiber_clock()
if wait_timeout < 0 then
Expand Down Expand Up @@ -246,10 +334,11 @@ function call.single(vshard_router, bucket_id, func_name, func_args, opts)
local timeout = opts.timeout or const.DEFAULT_VSHARD_CALL_TIMEOUT
local request_timeout = opts.mode == 'read' and opts.request_timeout or nil

local res, err = call_with_retry_and_recovery(vshard_router, replicaset, vshard_call_name,
func_name, func_args, {timeout = timeout, request_timeout = request_timeout}, true)
local res, err, replicaset_id = call_single_with_recovery(vshard_router, replicaset, vshard_call_name,
func_name, func_args, {timeout = timeout, request_timeout = request_timeout})

if err ~= nil then
return nil, wrap_vshard_err(vshard_router, err, func_name, nil, bucket_id)
return nil, wrap_vshard_err(err, func_name, replicaset_id)
end

if res == box.NULL then
Expand All @@ -276,16 +365,16 @@ function call.any(vshard_router, func_name, func_args, opts)

local deadline = fiber_clock() + timeout

for replicaset_id, replicaset in pairs(replicasets) do
for _, replicaset in pairs(replicasets) do
local wait_timeout = deadline - fiber_clock()

local is_timeout = wait_timeout < 0
if is_timeout then
wait_timeout = 0
end

local res, err = call_with_retry_and_recovery(vshard_router, replicaset, 'callro',
func_name, func_args, {timeout = wait_timeout}, false)
local res, err, replicaset_id = call_single_with_recovery(vshard_router, replicaset, 'callro',
func_name, func_args, {timeout = wait_timeout})

if err == nil then
if res == box.NULL then
Expand All @@ -302,7 +391,7 @@ function call.any(vshard_router, func_name, func_args, opts)
end
end

return nil, wrap_vshard_err(vshard_router, last_err, func_name, last_replicaset_id)
return nil, wrap_vshard_err(last_err, func_name, last_replicaset_id)
end

return call
5 changes: 2 additions & 3 deletions crud/common/map_call_cases/base_postprocessor.lua
Original file line number Diff line number Diff line change
Expand Up @@ -7,12 +7,11 @@ local BasePostprocessor = {}
-- @function new
--
-- @return[1] table postprocessor
function BasePostprocessor:new(vshard_router)
function BasePostprocessor:new()
local postprocessor = {
results = {},
early_exit = false,
errs = nil,
vshard_router = vshard_router,
storage_info = {},
}

Expand Down Expand Up @@ -68,7 +67,7 @@ function BasePostprocessor:collect(result_info, err_info)

if err ~= nil then
self.results = nil
self.errs = err_info.err_wrapper(self.vshard_router, err, unpack(err_info.wrapper_args))
self.errs = err_info.err_wrapper(err, unpack(err_info.wrapper_args))
self.early_exit = true

return self.early_exit
Expand Down
2 changes: 1 addition & 1 deletion crud/common/map_call_cases/batch_postprocessor.lua
Original file line number Diff line number Diff line change
Expand Up @@ -55,7 +55,7 @@ function BatchPostprocessor:collect(result_info, err_info)
err_to_wrap = err.err
end

local err_obj = err_info.err_wrapper(self.vshard_router, err_to_wrap, unpack(err_info.wrapper_args))
Comment thread
Serpentian marked this conversation as resolved.
local err_obj = err_info.err_wrapper(err_to_wrap, unpack(err_info.wrapper_args))
err_obj.operation_data = err.operation_data
err_obj.space_schema_hash = err.space_schema_hash

Expand Down
9 changes: 6 additions & 3 deletions crud/common/rebalance.lua
Original file line number Diff line number Diff line change
Expand Up @@ -25,10 +25,13 @@ local function safe_mode_bucket_trigger(_, new, space, op)
if space ~= '_bucket' then
return
end
-- We are interested only in two operations that indicate the beginning of bucket migration:
-- * We are receiving a bucket (new bucket with status RECEIVING)
-- * We are sending a bucket to another node (existing bucket status changes to SENDING)
-- Rebalancing start markers:
Comment thread
Serpentian marked this conversation as resolved.
-- * INSERT with RECEIVING - the bucket is being received.
-- * UPDATE to READONLY - since vshard 0.1.41 a transfer starts in
-- READONLY.
-- * REPLACE to SENDING - compatibility with older vshard versions.
if (op == 'INSERT' and new.status == vshard_consts.BUCKET.RECEIVING) or
(op == 'UPDATE' and new.status == vshard_consts.BUCKET.READONLY) or
(op == 'REPLACE' and new.status == vshard_consts.BUCKET.SENDING) then
local stored_safe_mode = schema.settings_space:get{ SAFE_MODE_STATUS }
if not stored_safe_mode or not stored_safe_mode.value then
Expand Down
5 changes: 3 additions & 2 deletions crud/common/sharding/init.lua
Original file line number Diff line number Diff line change
Expand Up @@ -18,7 +18,8 @@ function sharding.get_replicasets_by_bucket_id(vshard_router, bucket_id)
return nil, GetReplicasetsError:new("Failed to get replicaset for bucket_id %s: %s", bucket_id, err.err)
end

local replicaset_id = utils.get_replicaset_id(vshard_router, replicaset)
local replicaset_id = replicaset.id

if replicaset_id == nil then
return nil, GetReplicasetsError:new("Failed to get replicaset id for bucket_id %s replicaset", bucket_id)
end
Expand Down Expand Up @@ -314,7 +315,7 @@ function sharding.split_tuples_by_replicaset(vshard_router, tuples, space, opts)
sharding_data.bucket_id, err.err)
end

local replicaset_id = utils.get_replicaset_id(vshard_router, replicaset)
Comment thread
Serpentian marked this conversation as resolved.
local replicaset_id = replicaset.id
if replicaset_id == nil then
return nil, GetReplicasetsError:new(
"Failed to get replicaset id for bucket_id %s replicaset",
Expand Down
13 changes: 0 additions & 13 deletions crud/common/vshard_utils.lua
Original file line number Diff line number Diff line change
Expand Up @@ -100,19 +100,6 @@ function vshard_utils.get_self_vshard_replica_id()
end
end

function vshard_utils.get_replicaset_id(vshard_router, replicaset)
-- https://github.com/tarantool/vshard/issues/460.
local known_replicasets = vshard_router:routeall()

for known_replicaset_id, known_replicaset in pairs(known_replicasets) do
if known_replicaset == replicaset then
return known_replicaset_id
end
end

return nil
end

function vshard_utils.get_vshard_identification_mode()
-- https://github.com/tarantool/vshard/issues/460.
assert(vshard.storage.internal.current_cfg ~= nil, 'available only on vshard storage')
Expand Down
2 changes: 1 addition & 1 deletion crud/insert_many.lua
Original file line number Diff line number Diff line change
Expand Up @@ -214,7 +214,7 @@ local function call_insert_many_on_router(vshard_router, space_name, original_tu
return nil, {err}, const.NEED_SCHEMA_RELOAD
end

local postprocessor = BatchPostprocessor:new(vshard_router)
local postprocessor = BatchPostprocessor:new()

local rows, errs, storages_info = call.map(vshard_router, CRUD_INSERT_MANY_FUNC_NAME, nil, {
timeout = opts.timeout,
Expand Down
2 changes: 1 addition & 1 deletion crud/replace_many.lua
Original file line number Diff line number Diff line change
Expand Up @@ -214,7 +214,7 @@ local function call_replace_many_on_router(vshard_router, space_name, original_t
return nil, {err}, const.NEED_SCHEMA_RELOAD
end

local postprocessor = BatchPostprocessor:new(vshard_router)
local postprocessor = BatchPostprocessor:new()

local rows, errs, storages_info = call.map(vshard_router, CRUD_REPLACE_MANY_FUNC_NAME, nil, {
timeout = opts.timeout,
Expand Down
Loading
Loading