Skip to content
Merged
Show file tree
Hide file tree
Changes from 8 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
6 changes: 3 additions & 3 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -202,9 +202,9 @@ docker run --rm \
-listen :6479 \
-primary redis.internal:6379 \
-secondary elastickv.internal:6380 \
-elastickv-pool-size 4 \
-secondary-write-concurrency 2 \
-secondary-script-concurrency 1 \
-elastickv-pool-size 192 \
-secondary-write-concurrency 96 \
-secondary-script-concurrency 3 \
-mode dual-write
```

Expand Down
4 changes: 4 additions & 0 deletions adapter/redis.go
Original file line number Diff line number Diff line change
Expand Up @@ -116,6 +116,10 @@ const (

const (
redisDispatchTimeout = 10 * time.Second
// redisLuaDispatchTimeout gives EVAL/EVALSHA enough room for migration
// scripts that expand into thousands of Redis calls while keeping regular
// commands on the tighter dispatch deadline.
redisLuaDispatchTimeout = 30 * time.Second
// defaultRedisBlockWaitFallback is the safety-net poll interval for
// blocking-command wait loops when no in-process write signal arrives.
// Signals cover normal XADD / ZADD / ZINCRBY wakeups immediately; this
Expand Down
106 changes: 106 additions & 0 deletions adapter/redis_compat_helpers.go
Original file line number Diff line number Diff line change
Expand Up @@ -1112,6 +1112,112 @@ func (r *RedisServer) deleteLogicalKeyElems(ctx context.Context, key []byte, rea
return elems, existed, nil
}

func (r *RedisServer) deleteLogicalKeyElemsForType(ctx context.Context, key []byte, readTS uint64, typ redisValueType) ([]*kv.Elem[kv.OP], bool, error) {
switch typ {
case redisTypeNone:
return nil, false, nil
case redisTypeString:
elems, err := r.deleteStringLikeElems(ctx, key, readTS)
return elems, true, err
case redisTypeList:
return r.deleteListLogicalKeyElems(ctx, key, readTS)
case redisTypeHash:
return r.deleteHashLogicalKeyElems(ctx, key, readTS)
case redisTypeSet:
return r.deleteSetLogicalKeyElems(ctx, key, readTS)
case redisTypeZSet:
return r.deleteZSetLogicalKeyElems(ctx, key, readTS)
case redisTypeStream:
return r.deleteStreamLogicalKeyElems(ctx, key, readTS)
}
return nil, false, errors.WithStack(errors.AssertionFailedf("unknown redis type %v", typ))
}

func (r *RedisServer) deleteStringLikeElems(ctx context.Context, key []byte, readTS uint64) ([]*kv.Elem[kv.OP], error) {
var elems []*kv.Elem[kv.OP]
for _, internalKey := range [][]byte{
redisStrKey(key),
key, // legacy bare string key
redisHLLKey(key),
redisTTLKey(key),
} {
ok, err := r.store.ExistsAt(ctx, internalKey, readTS)
if err != nil {
return nil, errors.WithStack(err)
}
if ok {
elems = append(elems, &kv.Elem[kv.OP]{Op: kv.Del, Key: internalKey})
}
}
return elems, nil
}

func (r *RedisServer) deleteListLogicalKeyElems(ctx context.Context, key []byte, readTS uint64) ([]*kv.Elem[kv.OP], bool, error) {
elems, err := r.deleteStringLikeElems(ctx, key, readTS)
if err != nil {
return nil, false, err
}
listElems, err := r.deleteListElems(ctx, key, readTS)
if err != nil {
return nil, false, err
}
return append(elems, listElems...), true, nil
}

func (r *RedisServer) deleteHashLogicalKeyElems(ctx context.Context, key []byte, readTS uint64) ([]*kv.Elem[kv.OP], bool, error) {
elems, err := r.deleteStringLikeElems(ctx, key, readTS)
if err != nil {
return nil, false, err
}
elems = append(elems, &kv.Elem[kv.OP]{Op: kv.Del, Key: redisHashKey(key)})
hashElems, err := r.deleteWideColumnElems(ctx, readTS,
store.HashFieldScanPrefix(key), store.HashMetaKey(key), store.HashMetaDeltaScanPrefix(key))
if err != nil {
return nil, false, err
}
return append(elems, hashElems...), true, nil
}

func (r *RedisServer) deleteSetLogicalKeyElems(ctx context.Context, key []byte, readTS uint64) ([]*kv.Elem[kv.OP], bool, error) {
elems, err := r.deleteStringLikeElems(ctx, key, readTS)
if err != nil {
return nil, false, err
}
elems = append(elems, &kv.Elem[kv.OP]{Op: kv.Del, Key: redisSetKey(key)})
setElems, err := r.deleteWideColumnElems(ctx, readTS,
store.SetMemberScanPrefix(key), store.SetMetaKey(key), store.SetMetaDeltaScanPrefix(key))
if err != nil {
return nil, false, err
}
return append(elems, setElems...), true, nil
}

func (r *RedisServer) deleteZSetLogicalKeyElems(ctx context.Context, key []byte, readTS uint64) ([]*kv.Elem[kv.OP], bool, error) {
elems, err := r.deleteStringLikeElems(ctx, key, readTS)
if err != nil {
return nil, false, err
}
elems = append(elems, &kv.Elem[kv.OP]{Op: kv.Del, Key: redisZSetKey(key)})
zsetElems, err := r.deleteZSetWideColumnElems(ctx, key, readTS)
if err != nil {
return nil, false, err
}
return append(elems, zsetElems...), true, nil
}

func (r *RedisServer) deleteStreamLogicalKeyElems(ctx context.Context, key []byte, readTS uint64) ([]*kv.Elem[kv.OP], bool, error) {
elems, err := r.deleteStringLikeElems(ctx, key, readTS)
if err != nil {
return nil, false, err
}
elems = append(elems, &kv.Elem[kv.OP]{Op: kv.Del, Key: redisStreamKey(key)})
streamElems, err := r.deleteStreamWideColumnElems(ctx, key, readTS)
if err != nil {
return nil, false, err
}
return append(elems, streamElems...), true, nil
}

// deleteStreamWideColumnElems returns delete operations for all stream
// wide-column keys: the meta key (if it exists) and every entry under the
// entry scan prefix. Total results are capped at maxWideColumnItems to
Expand Down
2 changes: 1 addition & 1 deletion adapter/redis_lua.go
Original file line number Diff line number Diff line change
Expand Up @@ -115,7 +115,7 @@ func (r *RedisServer) runLuaScript(conn redcon.Conn, script string, evalArgs [][
return
}

ctx, cancel := context.WithTimeout(r.handlerContext(), redisDispatchTimeout)
ctx, cancel := context.WithTimeout(r.handlerContext(), redisLuaDispatchTimeout)
defer cancel()

start := time.Now()
Expand Down
19 changes: 14 additions & 5 deletions adapter/redis_lua_context.go
Original file line number Diff line number Diff line change
Expand Up @@ -3534,7 +3534,7 @@ func (c *luaScriptContext) commit() error {
}
sort.Strings(keys)

ctx, cancel := context.WithTimeout(c.scriptCtx(), redisDispatchTimeout)
ctx, cancel := context.WithTimeout(c.scriptCtx(), redisLuaDispatchTimeout)
defer cancel()

// Pre-allocate a commitTS so Delta key bytes can embed it before dispatch.
Expand Down Expand Up @@ -3658,18 +3658,27 @@ func (c *luaScriptContext) commitPlanForKey(ctx context.Context, key string, com
return luaKeyPlan{}, err
}

startType, err := c.server.keyTypeAt(ctx, []byte(key), c.startTS)
keyBytes := []byte(key)
rawStartType, err := c.server.rawKeyTypeAt(ctx, keyBytes, c.startTS)
if err != nil {
return luaKeyPlan{}, err
}
startType, err := c.server.applyTTLFilter(ctx, keyBytes, c.startTS, rawStartType)
if err != nil {
return luaKeyPlan{}, err
}
var deleteElems []*kv.Elem[kv.OP]
readKeys := luaWideFenceReadKeysForPlan([]byte(key), finalType, startType, valuePlan.preserveExisting)
readKeys := luaWideFenceReadKeysForPlan(keyBytes, finalType, startType, valuePlan.preserveExisting)
if !valuePlan.preserveExisting {
deleteElems, _, err = c.server.deleteLogicalKeyElems(ctx, []byte(key), c.startTS)
if c.everDeleted[key] && rawStartType != redisTypeNone {
deleteElems, _, err = c.server.deleteLogicalKeyElems(ctx, keyBytes, c.startTS)
} else {
deleteElems, _, err = c.server.deleteLogicalKeyElemsForType(ctx, keyBytes, c.startTS, rawStartType)
}
if err != nil {
return luaKeyPlan{}, err
}
deleteElems = append(deleteElems, redisTxnWideCollectionFenceElems([]byte(key))...)
deleteElems = append(deleteElems, redisTxnWideCollectionFenceElems(keyBytes)...)
}

dataElems := make([]*kv.Elem[kv.OP], 0, len(deleteElems)+len(valuePlan.elems))
Expand Down
7 changes: 4 additions & 3 deletions adapter/redis_peer_limiter.go
Original file line number Diff line number Diff line change
Expand Up @@ -10,9 +10,10 @@ import (

const (
redisPerPeerLimitEnv = "ELASTICKV_REDIS_PER_PEER_CONNECTIONS"
defaultRedisProxyPoolPeerCap = 64
defaultRedisDedicatedPeerHeadroom = 64
defaultRedisPerPeerConnectionCap = defaultRedisProxyPoolPeerCap + defaultRedisDedicatedPeerHeadroom
defaultRedisProxyPoolPeerCap = 192
defaultRedisProxyReplicasPerPeer = 2
defaultRedisDedicatedPeerHeadroom = 128
defaultRedisPerPeerConnectionCap = defaultRedisProxyPoolPeerCap*defaultRedisProxyReplicasPerPeer + defaultRedisDedicatedPeerHeadroom
redisPeerLimitError = "ERR max connections per client exceeded"
unknownRedisPeer = "unknown"
)
Expand Down
6 changes: 5 additions & 1 deletion adapter/redis_peer_limiter_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -15,7 +15,11 @@ func TestRedisPeerLimiterDefaultMatchesProxyPool(t *testing.T) {
t.Setenv(redisPerPeerLimitEnv, "")
limiter := newDefaultRedisPeerLimiter()
require.NotNil(t, limiter)
require.Equal(t, proxy.DefaultElasticKVBackendOptions().PoolSize+defaultRedisDedicatedPeerHeadroom, limiter.limit)
require.Equal(t, proxy.DefaultElasticKVBackendOptions().PoolSize, defaultRedisProxyPoolPeerCap)
require.Equal(t,
proxy.DefaultElasticKVBackendOptions().PoolSize*defaultRedisProxyReplicasPerPeer+defaultRedisDedicatedPeerHeadroom,
limiter.limit,
)
}

func TestRedisPeerLimiterRejectsAndReleases(t *testing.T) {
Expand Down
104 changes: 104 additions & 0 deletions adapter/redis_txn_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -639,6 +639,110 @@ func TestLuaWideFenceReadKeysForPlan(t *testing.T) {
require.Nil(t, luaWideFenceReadKeysForPlan(key, redisTypeString, redisTypeString, true))
}

type luaCleanupScanTrackingStore struct {
store.MVCCStore
fullScanStarts [][]byte
}

func (s *luaCleanupScanTrackingStore) ScanAt(ctx context.Context, start []byte, end []byte, limit int, ts uint64) ([]*store.KVPair, error) {
if limit == store.MaxDeltaScanLimit {
s.fullScanStarts = append(s.fullScanStarts, bytes.Clone(start))
}
return s.MVCCStore.ScanAt(ctx, start, end, limit, ts)
}

func newLuaCommitPlanTestContext(server *RedisServer, startTS uint64) *luaScriptContext {
return &luaScriptContext{
server: server,
startTS: startTS,
touched: map[string]struct{}{},
readKeys: map[string][]byte{},
deleted: map[string]bool{},
everDeleted: map[string]bool{},
negativeType: map[string]bool{},
strings: map[string]*luaStringState{},
lists: map[string]*luaListState{},
hashes: map[string]*luaHashState{},
sets: map[string]*luaSetState{},
zsets: map[string]*luaZSetState{},
streams: map[string]*luaStreamState{},
ttls: map[string]*luaTTLState{},
}
}

func TestLuaCommitPlanForAbsentRewriteSkipsFullLogicalCleanupScans(t *testing.T) {
t.Parallel()

ctx := context.Background()
base := store.NewMVCCStore()
tracking := &luaCleanupScanTrackingStore{MVCCStore: base}
server := NewRedisServer(nil, "", tracking, newLocalAdapterCoordinator(base), nil, nil)
key := "lua:absent-rewrite"

scriptCtx := newLuaCommitPlanTestContext(server, 10)
scriptCtx.strings[key] = &luaStringState{loaded: true, exists: true, dirty: true, value: []byte("v")}
scriptCtx.ttls[key] = &luaTTLState{loaded: true}

plan, err := scriptCtx.commitPlanForKey(ctx, key, 11)
require.NoError(t, err)
require.Empty(t, tracking.fullScanStarts)
require.True(t, elemKeysContain(plan.elems, redisStrKey([]byte(key))))
require.True(t, elemKeysContain(plan.elems, redisTxnWideHashFenceKey([]byte(key))))
}

func TestLuaCommitPlanForExistingListRewriteOnlyScansListCleanup(t *testing.T) {
t.Parallel()

ctx := context.Background()
base := store.NewMVCCStore()
tracking := &luaCleanupScanTrackingStore{MVCCStore: base}
server := NewRedisServer(nil, "", tracking, newLocalAdapterCoordinator(base), nil, nil)
key := []byte("lua:list-rewrite")
keyString := string(key)
meta, err := store.MarshalListMeta(store.ListMeta{Len: 1})
require.NoError(t, err)
require.NoError(t, base.PutAt(ctx, store.ListMetaKey(key), meta, 10, 0))
require.NoError(t, base.PutAt(ctx, listItemKey(key, 0), []byte("old"), 10, 0))

scriptCtx := newLuaCommitPlanTestContext(server, 11)
scriptCtx.lists[keyString] = &luaListState{
loaded: true,
exists: true,
dirty: true,
materialized: true,
values: []string{"new"},
}
scriptCtx.ttls[keyString] = &luaTTLState{loaded: true}

_, err = scriptCtx.commitPlanForKey(ctx, keyString, 12)
require.NoError(t, err)
requireScanStartsIncludePrefix(t, tracking.fullScanStarts, append(append([]byte(nil), []byte(store.ListItemPrefix)...), key...))
requireScanStartsIncludePrefix(t, tracking.fullScanStarts, store.ListMetaDeltaScanPrefix(key))
requireScanStartsIncludePrefix(t, tracking.fullScanStarts, store.ListClaimScanPrefix(key))
requireScanStartsExcludePrefix(t, tracking.fullScanStarts, store.HashFieldScanPrefix(key))
requireScanStartsExcludePrefix(t, tracking.fullScanStarts, store.SetMemberScanPrefix(key))
requireScanStartsExcludePrefix(t, tracking.fullScanStarts, store.ZSetMemberScanPrefix(key))
requireScanStartsExcludePrefix(t, tracking.fullScanStarts, store.ZSetScoreScanPrefix(key))
requireScanStartsExcludePrefix(t, tracking.fullScanStarts, store.StreamEntryScanPrefix(key))
}

func requireScanStartsIncludePrefix(t *testing.T, starts [][]byte, prefix []byte) {
t.Helper()
for _, start := range starts {
if bytes.HasPrefix(start, prefix) {
return
}
}
t.Fatalf("expected a scan under prefix %q, got %q", prefix, starts)
}

func requireScanStartsExcludePrefix(t *testing.T, starts [][]byte, prefix []byte) {
t.Helper()
for _, start := range starts {
require.Falsef(t, bytes.HasPrefix(start, prefix), "unexpected scan under prefix %q in %q", prefix, starts)
}
}

func TestRedisTxnSetReplacementConflictsWithConcurrentWideHashWrite(t *testing.T) {
t.Parallel()

Expand Down
19 changes: 13 additions & 6 deletions cmd/redis-proxy/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -19,9 +19,11 @@ import (
)

const (
sentryFlushTimeout = 2 * time.Second
metricsShutdownTimeout = 5 * time.Second
secondaryConcurrencyDivisor = 2
sentryFlushTimeout = 2 * time.Second
metricsShutdownTimeout = 5 * time.Second
secondaryWriteConcurrencyDivisor = 2
secondaryScriptConcurrencyDivisor = 32
secondaryScriptConcurrencyCap = 3
)

func main() {
Expand Down Expand Up @@ -53,13 +55,14 @@ func run() error {
flag.IntVar(&primaryPoolSize, "primary-pool-size", primaryPoolSize, "Primary Redis backend connection pool size")
flag.IntVar(&elasticKVPoolSize, "elastickv-pool-size", elasticKVPoolSize, "ElasticKV backend connection pool size")
flag.IntVar(&secondaryWriteConcurrency, "secondary-write-concurrency", secondaryWriteConcurrency, "Maximum concurrent asynchronous secondary writes including scripts (0 = half of secondary backend pool size)")
flag.IntVar(&secondaryScriptConcurrency, "secondary-script-concurrency", secondaryScriptConcurrency, "Maximum concurrent asynchronous secondary Lua-script writes within the write limit (0 = half of secondary write concurrency)")
flag.IntVar(&secondaryScriptConcurrency, "secondary-script-concurrency", secondaryScriptConcurrency, "Maximum concurrent asynchronous secondary Lua-script writes within the write limit (0 = secondary write concurrency / 32, minimum 1, capped at 3)")
flag.IntVar(&secondaryBlockingReplayConcurrency, "secondary-blocking-replay-concurrency", secondaryBlockingReplayConcurrency, "Maximum concurrent asynchronous secondary mutating blocking-command replays (0 = capped remaining secondary backend pool capacity after writes)")
flag.IntVar(&secondaryWriteQueueSize, "secondary-write-queue-size", secondaryWriteQueueSize, "Maximum queued asynchronous secondary writes (0 = derived from write concurrency)")
flag.IntVar(&secondaryScriptQueueSize, "secondary-script-queue-size", secondaryScriptQueueSize, "Maximum queued asynchronous secondary Lua-script writes (0 = derived from script concurrency)")
flag.IntVar(&secondaryBlockingReplayQueueSize, "secondary-blocking-replay-queue-size", secondaryBlockingReplayQueueSize, "Maximum queued asynchronous secondary mutating blocking-command replays (0 = derived from blocking replay concurrency)")
flag.StringVar(&modeStr, "mode", "dual-write", "Proxy mode: redis-only, dual-write, dual-write-shadow, elastickv-primary, elastickv-only")
flag.DurationVar(&cfg.SecondaryTimeout, "secondary-timeout", cfg.SecondaryTimeout, "Secondary write timeout")
flag.DurationVar(&cfg.SecondaryScriptTimeout, "secondary-script-timeout", cfg.SecondaryScriptTimeout, "Secondary Lua-script write timeout (0 = secondary-timeout)")
Comment thread
coderabbitai[bot] marked this conversation as resolved.
Outdated

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge Extend Redis secondary deadlines for primary-cutover scripts

In elastickv-primary mode the async secondary for script replays is Redis (primaryOpts), whose backend options still keep the 3s ReadTimeout; this new -secondary-script-timeout only changes the task context. A Lua script that succeeds on ElasticKV but takes longer than 3s to replay to Redis will still hit go-redis' socket deadline and be recorded as a secondary write failure despite the advertised 5m default, so the Redis backend needs a secondary-specific read timeout when it is used for async script replay.

Useful? React with 👍 / 👎.

flag.DurationVar(&cfg.ShadowTimeout, "shadow-timeout", cfg.ShadowTimeout, "Shadow read timeout")
flag.StringVar(&cfg.SentryDSN, "sentry-dsn", cfg.SentryDSN, "Sentry DSN (empty = disabled)")
flag.StringVar(&cfg.SentryEnv, "sentry-env", cfg.SentryEnv, "Sentry environment")
Expand Down Expand Up @@ -283,11 +286,15 @@ func secondaryBackendPoolSize(mode proxy.ProxyMode, primaryPoolSize, elasticKVPo
}

func defaultSecondaryWriteConcurrency(poolSize int) int {
return atLeastOne(poolSize / secondaryConcurrencyDivisor)
return atLeastOne(poolSize / secondaryWriteConcurrencyDivisor)
}

func defaultSecondaryScriptConcurrency(writeConcurrency int) int {
return atLeastOne(writeConcurrency / secondaryConcurrencyDivisor)
concurrency := atLeastOne(writeConcurrency / secondaryScriptConcurrencyDivisor)
if concurrency > secondaryScriptConcurrencyCap {
return secondaryScriptConcurrencyCap
}
return concurrency
}

func defaultSecondaryBlockingReplayConcurrency(poolSize, writeConcurrency int) int {
Expand Down
Loading
Loading