Visitar URL original
Add cross-shard pipelines and transactional pipeline blocks to RedisCluster by ebubekiryigit · Pull Request #2798 · phpredis/phpredis · GitHub
Skip to content

Add cross-shard pipelines and transactional pipeline blocks to RedisCluster - #2798

Open
ebubekiryigit wants to merge 21 commits into
phpredis:developfrom
ebubekiryigit:feature/redis-cluster-pipeline-support
Open

ebubekiryigit wants to merge 21 commits into
phpredis:developfrom
ebubekiryigit:feature/redis-cluster-pipeline-support

Conversation

@ebubekiryigit

@ebubekiryigit ebubekiryigit commented Jan 21, 2026 •

Copy link
Copy Markdown
Contributor

Summary

This PR adds native pipelining to RedisCluster while keeping the public API and mode transitions aligned with standalone Redis:

  • pipeline() and multi(Redis::PIPELINE) start the same non-atomic pipeline mode.
  • Non-atomic pipelines may contain commands for different hash slots and shards.
  • pipeline()->multi() creates a Redis transaction inside the pipeline. Each such transaction is restricted to one hash slot.
  • Existing multi() behavior is preserved: PhpRedis may coordinate transactions on multiple cluster nodes, with atomicity provided independently by each node rather than globally across shards.
  • Distributed MGET, MSET, MSETNX, DEL, and UNLINK operations 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

API Mode Slot/shard behavior Atomicity
Normal commands ATOMIC Existing RedisCluster routing; supported distributed commands may fan out across shards Commands execute immediately; a distributed helper is not a cross-shard transaction
pipeline() PIPELINE Cross-slot and cross-shard commands are allowed Non-atomic; partial execution is possible if a node fails
multi(Redis::PIPELINE) PIPELINE Identical to pipeline() Non-atomic
pipeline()->multi() PIPELINE | MULTI Every MULTI block must stay in one hash slot; outer pipeline commands and separate MULTI blocks may use other shards Atomic only within that single-slot Redis transaction
multi() MULTI Existing multi-node behavior remains supported Atomic per participating Redis node, not globally across shards

For example:

$pipe = $redis->pipeline();
$pipe->set('{a}one', '1')
     ->set('{b}two', '2');
$replies = $pipe->exec();

multi(Redis::PIPELINE) is an alias for the same non-atomic mode:

$pipe = $redis->multi(Redis::PIPELINE);
$pipe->set('{a}one', '1')
     ->get('{a}one');
$replies = $pipe->exec();

A transaction inside a pipeline uses the standalone PhpRedis nesting model. The inner exec() closes and queues the MULTI block; the outer exec() sends the pipeline:

$pipe = $redis->pipeline();
$pipe->multi()
     ->set('{account:1}balance', '100')
     ->get('{account:1}balance')
     ->exec();

$replies = $pipe->exec(); // [[true, '100']]

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:

  • Queueing performs no socket I/O. Commands are buffered per resolved destination socket.
  • exec() writes one complete buffer per participating socket and rejects short/partial writes.
  • Every queued reply is bound to the exact socket that received the command. Reply processing does not look the socket up again through a potentially changed slot map.
  • Replies are folded into the original PHP command order even though commands are grouped by node for transport.
  • Normal Redis command errors are represented as false for that result while all remaining replies are consumed, leaving the socket clean.
  • On a write failure, truncated reply, invalid framing, 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.
  • Pipelines are not automatically replayed after an ambiguous failure. Some commands may already have executed, so retrying the whole batch could duplicate side effects.
  • Queued command contexts have explicit destructors and ownership-transfer state. exec(), discard(), error paths, and object destruction release each context exactly once.
  • Reentrant commands on the same RedisCluster object 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 directed WAIT/WAITAOF calls.

Distributed multi-key commands

For MGET, MSET, MSETNX, DEL, and UNLINK, the RedisCluster abstraction remains consistent with its existing immediate-command behavior:

  1. Keys are prefixed and hashed normally.
  2. The logical command is split into per-slot commands.
  3. Each fragment is added to the appropriate node buffer.
  4. Fragment replies are combined into one logical result in key/argument order.

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. MSETNX therefore 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:

  • Cross-slot and interleaved cross-node pipelines with stable reply ordering.
  • Equivalence of pipeline() and multi(Redis::PIPELINE).
  • Single-slot pipelined MULTI blocks, multiple blocks on different shards, empty blocks, and hash tags.
  • Prefix and serializer behavior, missing keys, distributed multi-key commands, and command errors.
  • Wrong-type replies and ACL-denied MSET, MSETNX, DEL, and UNLINK fragments without leaving unread responses.
  • WATCH aborts, rejected commands, invalid mode transitions, and reentrant unserialization callbacks.
  • Cleanup/leak-oriented lifecycle cases: distributed contexts followed by exec(), discard(), pipeline abort, and object destruction.
  • Regression tests that never enter pipeline mode, covering existing immediate commands, distributed multi(), response ordering, error handling, DISCARD, and subsequent socket reuse.
  • A non-pipeline regression test covers MOVED followed by a partial slot map, with persistent connections enabled and disabled, and verifies subsequent reads from covered slots.
  • Failed-socket reconnect recovery through both pipeline entry APIs, with and without persistent connections, and nested EXEC return-type validation.

Additional validation:

  • Full RedisCluster suite, including serializer, compression, and session coverage: passed.
  • Full standalone Redis suite: passed.
  • Pipeline-disabled regression tests were run against both this branch and a separately built origin/develop module: identical results, all passed.
  • Manual resharding/redirection, node failure, and CLUSTERDOWN scenarios were exercised to verify failure returns, state reset, and clean socket recovery.
  • Build and generated arginfo consistency checks, git diff --check, and repeated post-failure commands: passed.

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.
@ebubekiryigit

Copy link
Copy Markdown
Contributor Author

This change also addresses the long-standing discussion in: #1910

@michael-grunder

Copy link
Copy Markdown
Member

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.

@ebubekiryigit

Copy link
Copy Markdown
Contributor Author

Hey @michael-grunder,

Thanks for the review! Let me know if there’s anything I should clarify or adjust.

@aphofstede

Copy link
Copy Markdown

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.
@ebubekiryigit ebubekiryigit changed the title feat(redis-cluster): implemented single-slot pipeline support and tests Add cross-shard pipelines and transactional pipeline blocks to RedisCluster Sep 4, 2026
@ebubekiryigit
ebubekiryigit marked this pull request as draft September 4, 2026 10:37
@ebubekiryigit
ebubekiryigit marked this pull request as ready for review September 7, 2026 05:45
@ebubekiryigit

Copy link
Copy Markdown
Contributor Author

The implementation in 134013d was limited to single-slot pipelines. This revision expands RedisCluster pipelining across slots and cluster nodes while keeping the behavior aligned with standalone Redis:

  • pipeline() and multi(Redis::PIPELINE) create a non-atomic pipeline.
  • pipeline()->multi() creates an atomic, single-slot transaction.
  • The existing RedisCluster multi() behavior remains unchanged.

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, WATCH/UNWATCH remain unsupported in pipeline mode, and pipelines are not replayed after ambiguous failures. Parallel dispatch to cluster nodes and pipeline-aware WATCH support are intentionally left for future work.

I’m looking forward to your review and excited to see this feature move forward.

@ebubekiryigit

Copy link
Copy Markdown
Contributor Author

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 develop again until review, since it keeps creating conflicts while the PR is waiting.

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.
@ebubekiryigit
ebubekiryigit force-pushed the feature/redis-cluster-pipeline-support branch from fd54913 to 6104bf5 Compare October 4, 2026 12:00
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.
@michael-grunder michael-grunder self-assigned this Oct 5, 2026
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.

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants