From bb1879ba35bb0aace7c5f3c66b9bdf2ffb7a1e6c Mon Sep 17 00:00:00 2001 From: Wang Qi Date: Mon, 3 Aug 2026 15:44:59 +0800 Subject: [PATCH] REDIS: allkeys-lru -> volatile-lru to avoid got evict if lack of memory (#17720) Follow on PR: #17707 --- docker/docker-compose-base.yml | 2 +- internal/agent/canvas/run_tracker.go | 21 +++++++++++++++------ internal/service/chunk_types.go | 9 +++++---- internal/service/dataset/index.go | 2 +- 4 files changed, 22 insertions(+), 12 deletions(-) diff --git a/docker/docker-compose-base.yml b/docker/docker-compose-base.yml index 1afffe7bc7..286294f823 100644 --- a/docker/docker-compose-base.yml +++ b/docker/docker-compose-base.yml @@ -235,7 +235,7 @@ services: redis: # swr.cn-north-4.myhuaweicloud.com/ddn-k8s/docker.io/valkey/valkey:8 image: valkey/valkey:8 - command: ["redis-server", "--requirepass", "${REDIS_PASSWORD}", "--maxmemory", "128mb", "--maxmemory-policy", "allkeys-lru"] + command: ["redis-server", "--requirepass", "${REDIS_PASSWORD}", "--maxmemory", "128mb", "--maxmemory-policy", "volatile-lru"] env_file: .env ports: - ${REDIS_PORT}:6379 diff --git a/internal/agent/canvas/run_tracker.go b/internal/agent/canvas/run_tracker.go index bc138e6140..541e77a8b6 100644 --- a/internal/agent/canvas/run_tracker.go +++ b/internal/agent/canvas/run_tracker.go @@ -100,6 +100,15 @@ func NewRunTrackerWithClient(client *redis.Client, ttl time.Duration) *RunTracke return &RunTracker{client: client, ttl: ttl} } +func (t *RunTracker) setRunFields(ctx context.Context, runID string, values ...interface{}) error { + key := runKey(runID) + pipe := t.client.Pipeline() + pipe.HSet(ctx, key, values...) + pipe.Expire(ctx, key, t.ttl) + _, err := pipe.Exec(ctx) + return err +} + // Start records a new run as in-progress. canvasID and tenantID identify // the source DSL and tenant; parentRunID may be empty for fresh runs and // carries the source run-id for resume chains (R1 in plan ยง2.6). @@ -134,7 +143,7 @@ func (t *RunTracker) AttachCheckpoint(ctx context.Context, runID, checkpointID s if t == nil || t.client == nil { return errors.New("run tracker: redis client not initialized") } - return t.client.HSet(ctx, runKey(runID), "checkpoint_id", checkpointID).Err() + return t.setRunFields(ctx, runID, "checkpoint_id", checkpointID) } // AttachInterrupt persists the eino interrupt id that paused this run @@ -147,7 +156,7 @@ func (t *RunTracker) AttachInterrupt(ctx context.Context, runID, interruptID str if t == nil || t.client == nil { return errors.New("run tracker: redis client not initialized") } - return t.client.HSet(ctx, runKey(runID), runFieldInterruptID, interruptID).Err() + return t.setRunFields(ctx, runID, runFieldInterruptID, interruptID) } // GetInterruptID returns the persisted interrupt id for runID. The bool is @@ -186,10 +195,10 @@ func (t *RunTracker) MarkSucceeded(ctx context.Context, runID string) error { if t == nil || t.client == nil { return errors.New("run tracker: redis client not initialized") } - return t.client.HSet(ctx, runKey(runID), + return t.setRunFields(ctx, runID, "status", runStatusSucceeded, "finished_at", time.Now().UnixMilli(), - ).Err() + ) } // MarkFailed transitions the run to status=2 and records the reason. @@ -197,11 +206,11 @@ func (t *RunTracker) MarkFailed(ctx context.Context, runID, reason string) error if t == nil || t.client == nil { return errors.New("run tracker: redis client not initialized") } - return t.client.HSet(ctx, runKey(runID), + return t.setRunFields(ctx, runID, "status", runStatusFailed, "finished_at", time.Now().UnixMilli(), "failure_reason", reason, - ).Err() + ) } // MarkCancelled transitions the run to status=3 and sets the cancel flag. diff --git a/internal/service/chunk_types.go b/internal/service/chunk_types.go index 4f1e1c9817..625a121e00 100644 --- a/internal/service/chunk_types.go +++ b/internal/service/chunk_types.go @@ -21,14 +21,15 @@ import ( "encoding/json" "fmt" "math" - "ragflow/internal/common" - "ragflow/internal/engine/redis" - "ragflow/internal/entity" "strings" + "time" + "ragflow/internal/common" "ragflow/internal/dao" "ragflow/internal/engine" + "ragflow/internal/engine/redis" "ragflow/internal/engine/types" + "ragflow/internal/entity" "ragflow/internal/tokenizer" "ragflow/internal/utility" @@ -223,7 +224,7 @@ func (s *ChunkService) cancelAllTasksOfDoc(ctx context.Context, docID string) er if task == nil { continue } - redisClient.Set(ctx, fmt.Sprintf("%s-cancel", task.ID), "x", 0) + redisClient.Set(ctx, fmt.Sprintf("%s-cancel", task.ID), "x", time.Hour) } return nil diff --git a/internal/service/dataset/index.go b/internal/service/dataset/index.go index 5909e05bd8..8f0eda7ca7 100644 --- a/internal/service/dataset/index.go +++ b/internal/service/dataset/index.go @@ -690,7 +690,7 @@ func (d *DatasetService) DeleteIndex(ctx context.Context, userID, datasetID, ind if taskID != "" { redisClient := redisengine.Get() - if redisClient == nil || !redisClient.Set(ctx, fmt.Sprintf("%s-cancel", taskID), "x", 0) { + if redisClient == nil || !redisClient.Set(ctx, fmt.Sprintf("%s-cancel", taskID), "x", time.Hour) { common.Warn("Failed to set dataset index cancellation marker", zap.String("dataset_id", datasetID), zap.String("task_id", taskID)) } if err := dao.DB.Unscoped().Where("id = ?", taskID).Delete(&entity.Task{}).Error; err != nil {