Repository navigation
Add cross-shard pipelines and transactional pipeline blocks to RedisCluster - #2798
ebubekiryigit wants to merge 21 commits into
Conversation
Implement RedisCluster pipeline mode with slot enforcement and queue buffering, add pipeline exec handling, and reset logic for errors/discard. Update stubs, arginfo, and cluster docs. Expand RedisCluster tests to cover pipeline/multi edge cases, cross-slot errors, mode transitions, and invalid operations.
|
This change also addresses the long-standing discussion in: #1910 |
|
Hi thanks for the PR. I'm not opposed to adding the feature by any means. I'll play around with it locally this week. |
|
Hey @michael-grunder, Thanks for the review! Let me know if there’s anything I should clarify or adjust. |
|
Any chances of this getting integrated? |
…ter-pipeline-support # Conflicts: # cluster_library.h # redis_cluster.c
Align pipeline and MULTI state transitions with standalone Redis, preserve legacy distributed MULTI behavior, and harden socket/context cleanup across failure paths.
- distinguish valid false/null replies from decode failures - disconnect participants after malformed transaction responses - expand cross-node, lifecycle, and compatibility tests
|
The implementation in
Test coverage now includes reply ordering, distributed commands, mode transitions, malformed and truncated responses, failure recovery, socket cleanup, persistent connections, and sessions. The implementation has been validated with Redis 6.2, 7.4, 8.2, and 8.8. This revision also addresses the failure-handling, MULTI, and distributed multi-key concerns raised by @michael-grunder. To keep the scope and risk bounded, I’m looking forward to your review and excited to see this feature move forward. |
|
Hi @michael-grunder, I’ve updated the PR and addressed the remaining pipeline concerns. Could you please approve the current workflow run and let me know if there are any blockers left for review/merge? I’ll avoid syncing with |
Accept valid nil and timeout replies, and fail logical MGET results on chunk errors. Refresh routing before the next pipeline without replaying commands or replacing the original failure. Preserve legacy stream and ACL results and cover reply, cleanup, and recovery regressions.
fd54913 to
6104bf5
Compare
Clear stale slot pointers, retire old sockets, and authenticate cached seeds. Add regression coverage for each recovery path.
Check the remapped node before accessing its socket. Add a non-pipeline regression for partial slot maps, persistent connections, and subsequent covered-slot reads.
Deduplicate sends through buffer state and fold distributed errors with a boolean flag. Restore handlers unreachable from pipelines and report the correct mode in pipeline warnings.
Cover reply families, transport failures, topology recovery, cache refresh, persistent connections, permissions, and legacy MULTI cleanup. Assert abort cleanup before remapping, distinguish stale replies, verify packed bytes, and count master reads independently of replica lag.
Wait for seed slot maps to converge after MOVED/ASK slot changes. Count refreshes on replica seeds too, without retrying pipelines.
Declare the client return used by nested EXEC and regenerate arginfo. Keep getTransferredBytes unchanged and assert the fluent return and Reflection signature in the existing transaction test.
Detach each node buffer before reconnecting, so socket cleanup cannot free the pending write. Release it after both successful and failed sends. Cover recovery through both pipeline entry APIs, with persistent and non-persistent connections, using the existing reply-failure fixture.
Document the previous multi(PIPELINE) fallback, fatal entry from MULTI, and ATOMIC reset after an uncovered slot or transaction slot violation. Keep the changes confined to the new Pipelining section.
Summary
This PR adds native pipelining to
RedisClusterwhile keeping the public API and mode transitions aligned with standaloneRedis:pipeline()andmulti(Redis::PIPELINE)start the same non-atomic pipeline mode.pipeline()->multi()creates a Redis transaction inside the pipeline. Each such transaction is restricted to one hash slot.multi()behavior is preserved: PhpRedis may coordinate transactions on multiple cluster nodes, with atomicity provided independently by each node rather than globally across shards.MGET,MSET,MSETNX,DEL, andUNLINKoperations continue to be split by slot and are folded back into one logical result in the original command order.This supersedes the original single-slot-only implementation in this PR and addresses the design questions raised in #1910, particularly Michael's failure-handling, MULTI, and multi-key command concerns.
API and atomicity model
ATOMICpipeline()PIPELINEmulti(Redis::PIPELINE)PIPELINEpipeline()pipeline()->multi()PIPELINE | MULTImulti()MULTIFor example:
multi(Redis::PIPELINE)is an alias for the same non-atomic mode:A transaction inside a pipeline uses the standalone PhpRedis nesting model. The inner
exec()closes and queues the MULTI block; the outerexec()sends the pipeline:Keys within that MULTI block must resolve to the same slot, including after prefixing and hash-tag processing. This is a single Redis-node transaction; resharding may still make it fail rather than providing a cross-node atomicity guarantee.
Previously, multi(Redis::PIPELINE) warned and fell back to regular MULTI. It now starts a non-atomic pipeline, matching standalone Redis. Starting a pipeline from an active MULTI is a fatal error. An uncovered slot or a transaction slot violation discards the pending pipeline and restores ATOMIC mode; subsequent commands execute immediately.
Failure and socket-safety model
The implementation is intentionally conservative because a cross-shard pipeline can be partially delivered or partially executed:
exec()writes one complete buffer per participating socket and rejects short/partial writes.falsefor that result while all remaining replies are consumed, leaving the socket clean.MOVED/ASK, or another pipeline-level failure, all participating sockets are disconnected before state is reset. This prevents unread replies or a partially written request from contaminating a later command or a persistent connection.exec(),discard(), error paths, and object destruction release each context exactly once.RedisClusterobject are rejected while pipeline replies are being decoded.Failure recovery refreshes routing before the next pipeline starts, without replaying failed commands. At least one configured seed must remain reachable.
The shared remap fixes also protect legacy atomic commands: cached seeds retain authentication, rebuilt slot maps clear stale entries, and uncovered slots after MOVED raise an exception.
Commands that cannot be safely represented by this queue are rejected before any pipeline write. This includes directed-node commands,
KEYS/scan-style operations, Pub/Sub subscription operations,WATCH/UNWATCH, and directedWAIT/WAITAOFcalls.Distributed multi-key commands
For
MGET,MSET,MSETNX,DEL, andUNLINK, the RedisCluster abstraction remains consistent with its existing immediate-command behavior:In a non-atomic pipeline these commands may span shards. Inside a pipelined MULTI block they are accepted only when every key maps to the block's single slot.
MSETNXtherefore retains RedisCluster's existing per-slot result semantics when it is distributed; it does not become globally atomic across shards.Testing
The updated test coverage focuses on the failure modes and ownership concerns from #1910:
pipeline()andmulti(Redis::PIPELINE).MSET,MSETNX,DEL, andUNLINKfragments without leaving unread responses.WATCHaborts, rejected commands, invalid mode transitions, and reentrant unserialization callbacks.exec(),discard(), pipeline abort, and object destruction.multi(), response ordering, error handling,DISCARD, and subsequent socket reuse.Additional validation:
RedisClustersuite, including serializer, compression, and session coverage: passed.Redissuite: passed.origin/developmodule: identical results, all passed.CLUSTERDOWNscenarios were exercised to verify failure returns, state reset, and clean socket recovery.git diff --check, and repeated post-failure commands: passed.