技术分析报告 — 从 Redis Streams 到 Kafka 的完整改造方案
RAGFlow v0.25.6 使用 Redis (Valkey) Streams 作为异步任务消息代理。 这是一个基于 Redis Stream 消费者组的自研轻量级消息队列系统,无 Celery、RabbitMQ 或 Kafka 等第三方消息中间件。
单例 RedisDB 类封装了 Valkey Python 客户端 (v6.0.2),提供以下队列操作原语:
| 方法 | Redis 命令 | 功能 |
|---|---|---|
queue_product(queue, message) | XADD | 生产消息到 Stream,重试 3 次 |
queue_consumer(queue, group, consumer) | XREADGROUP | 消费消息,自动创建 Consumer Group |
get_unacked_iterator(queues, group, consumer) | 遍历 XPENDING | 重放未确认消息(启动恢复) |
RedisMsg.ack() | XACK | 确认消息消费完成 |
queue_info(queue, group) | XINFO GROUPS | 获取队列积压和消费者信息 |
get_pending_msg(queue, group) | XPENDING RANGE | 查看待处理消息列表 |
requeue_msg(queue, group, msg_id) | XRANGE + XADD + XACK | 重新入队消息 |
核心异步 Worker 进程,关键参数:
task_executor_{host_id}_{worker_no}MAX_CONCURRENT_TASKS=5 (asyncio Semaphore)| 函数 | 任务类型 | 触发场景 | 调用方 |
|---|---|---|---|
queue_tasks(doc, bucket, name, priority) |
parse | 用户上传文档到知识库 | Document API |
queue_dataflow(tenant_id, flow_id, task_id, ...) |
dataflow | Agent Canvas 工作流执行 | Canvas API |
queue_raptor_o_graphrag_tasks(doc, ty, ...) |
raptor / graphrag / mindmap | 用户运行 GraphRAG/Raptor/Mindmap | Knowledge Base API |
handle_save_to_memory_task |
memory | 记忆嵌入流水线 | Memory Message Service |
| 维度 | 现状 | 评级 |
|---|---|---|
| 消息持久化 | Redis Stream 持久化 (AOF/RDB),但 128MB 内存上限触发 LRU 淘汰 | ⚠️ 中等 |
| 高可用 | 单实例 Redis,无哨兵/集群 | ❌ 低 |
| 水平扩展 | Worker 可水平扩展,但所有 Worker 共享同一 Redis | ⚠️ 中等 |
| 消息回溯 | 支持未确认消息重放,但不支持按时间/offset 回溯 | ⚠️ 中等 |
| 死信队列 | 无 DLQ,失败消息永久 Pending | ❌ 低 |
| 消息顺序 | 支持,但无分区概念 | ✅ 良好 |
| 监控可观测 | 基础心跳 + Pending/Lag 计数,无消费延迟监控 | ⚠️ 中等 |
| 消息吞吐 | 适合中等规模 (受限于单 Redis 实例) | ⚠️ 中等 |
--maxmemory 128mb --maxmemory-policy allkeys-lru 意味着在高负载下,队列消息可能被 LRU 淘汰,导致任务丢失。生产环境不建议依赖此配置。
| 目标 | 当前 (Redis Streams) | 迁移后 (Kafka) |
|---|---|---|
| 消息持久化 | 内存优先,128MB 上限 | 磁盘持久化,TB 级存储,可配置保留策略 |
| 高可用 | 单点故障 | 多 Broker 集群,ISR 机制,自动故障转移 |
| 水平扩展 | Worker 可扩展,Broker 不可扩展 | Partition 级别并行,Broker 和 Consumer 均可横向扩展 |
| 消息回溯 | 仅 Pending 消息 | 基于 Offset 的任意时间点回溯 |
| 死信队列 | 无 | 原生支持 DLQ 模式 |
| 监控 | 基本心跳 | JMX 指标,Consumer Lag,与 Prometheus/Grafana 集成 |
| 生态集成 | 需自建 | Kafka Connect, ksqlDB, Schema Registry, Kafka Streams |
| Topic 名称 | 分区数 | 副本数 | 保留策略 | 用途 |
|---|---|---|---|---|
ragflow.task.parse |
8 | 3 | 7 天 | 文档解析/分块任务 (原 Priority 0) |
ragflow.task.parse.priority |
4 | 3 | 3 天 | 高优先级解析任务 (原 Priority 1) |
ragflow.task.dataflow |
4 | 3 | 7 天 | Agent Canvas 数据流任务 |
ragflow.task.graphrag |
2 | 3 | 7 天 | GraphRAG / Raptor / Mindmap 知识图谱任务 |
ragflow.task.memory |
2 | 3 | 3 天 | Memory 嵌入任务 |
ragflow.task.dlq |
2 | 3 | 30 天 | 死信队列 (所有 Topic 失败消息汇聚) |
tenant_id 的 hash 作为分区键,保证同一租户的任务有序处理。
key = f"{tenant_id}:{doc_id}:{task_id}" 或 key = f"{tenant_id}:{flow_id}"(dataflow)。
新建 rag/utils/message_queue.py,提供统一的 Producer/Consumer 接口,屏蔽底层实现差异:
# 新建: rag/utils/message_queue.py
from abc import ABC, abstractmethod
from typing import Any, Callable
from dataclasses import dataclass
import json
@dataclass
class QueueMessage:
"""统一消息模型"""
message_id: str
payload: dict
topic: str
partition: int = -1
offset: int = -1
timestamp: int = 0
def ack(self) -> bool:
"""确认消息消费完成"""
...
def nack(self, delay_ms: int = 0) -> bool:
"""拒绝消息 (可延迟重试)"""
...
class MessageQueueProducer(ABC):
"""消息生产者抽象"""
@abstractmethod
async def send(self, topic: str, key: str, value: dict) -> bool: ...
@abstractmethod
async def flush(self) -> None: ...
class MessageQueueConsumer(ABC):
"""消息消费者抽象"""
@abstractmethod
async def poll(self, timeout_ms: int = 5000) -> QueueMessage | None: ...
@abstractmethod
async def commit(self, message: QueueMessage) -> bool: ...
@abstractmethod
async def close(self) -> None: ...
class MessageQueueFactory:
"""工厂方法: 根据配置创建 Redis 或 Kafka 实现"""
_instance = None
@classmethod
def get_producer(cls) -> MessageQueueProducer:
if settings.MQ_BACKEND == "kafka":
return KafkaProducer(...)
return RedisStreamProducer(...)
@classmethod
def get_consumer(cls, group_id: str) -> MessageQueueConsumer:
if settings.MQ_BACKEND == "kafka":
return KafkaConsumer(group_id, ...)
return RedisStreamConsumer(group_id, ...)
所有现有 REDIS_CONN.queue_product(...) 调用替换为统一接口,影响函数:
| 文件 | 函数/位置 | 改动描述 |
|---|---|---|
| api/db/services/task_service.py:458 | queue_tasks() |
替换为 MQ_PRODUCER.send(topic, key, task_dict) |
| api/db/services/task_service.py:549 | queue_dataflow() |
替换为 MQ_PRODUCER.send("ragflow.task.dataflow", key, task_dict) |
| api/db/services/document_service.py:1101 | queue_raptor_o_graphrag_tasks() |
替换为 MQ_PRODUCER.send("ragflow.task.graphrag", key, task_dict) |
| api/db/joint_services/memory_message_service.py:414 | handle_save_to_memory_task (生产部分) |
替换为 MQ_PRODUCER.send("ragflow.task.memory", key, task_dict) |
# 新建: rag/utils/kafka_producer.py (简化示意)
import asyncio
from aiokafka import AIOKafkaProducer
from common import settings
class KafkaProducer(MessageQueueProducer):
def __init__(self):
self._producer = AIOKafkaProducer(
bootstrap_servers=settings.KAFKA_BOOTSTRAP_SERVERS,
key_serializer=lambda k: k.encode('utf-8'),
value_serializer=lambda v: json.dumps(v).encode('utf-8'),
acks='all', # 等待所有 ISR 副本确认
compression_type='gzip', # 压缩传输
max_in_flight_requests_per_connection=1, # 保证顺序
retries=5,
)
async def start(self):
await self._producer.start()
async def send(self, topic: str, key: str, value: dict) -> bool:
try:
result = await self._producer.send_and_wait(topic, key=key, value=value)
logging.debug(f"Kafka produce OK: {topic}[{result.partition}]@{result.offset}")
return True
except Exception as e:
logging.exception(f"Kafka produce failed: {e}")
return False
async def flush(self):
await self._producer.flush()
rag/svr/task_executor.py 的 collect() 函数从 Redis XREADGROUP 模式改造为 Kafka Consumer 模式:
# 改造前后对比 — task_executor.py collect()
# 改造前 (Redis Streams):
async def collect():
redis_msg = REDIS_CONN.queue_consumer(queue, group, consumer_name)
if redis_msg:
msg = redis_msg.get_message()
task = TaskService.get_task(msg["id"])
return redis_msg, task
# 改造后 (Kafka):
async def collect():
msg = await KAFKA_CONSUMER.poll(timeout_ms=5000)
if msg:
payload = msg.payload
task = TaskService.get_task(payload["id"])
return msg, task # QueueMessage 替代 RedisMsg
# 新建: rag/utils/kafka_consumer.py (简化示意)
import asyncio
from aiokafka import AIOKafkaConsumer
class KafkaConsumer(MessageQueueConsumer):
def __init__(self, group_id: str, topics: list[str]):
self._consumer = AIOKafkaConsumer(
*topics,
bootstrap_servers=settings.KAFKA_BOOTSTRAP_SERVERS,
group_id=group_id,
key_deserializer=lambda k: k.decode('utf-8'),
value_deserializer=lambda v: json.loads(v.decode('utf-8')),
enable_auto_commit=False, # 手动提交 → 等效 XACK
auto_offset_reset='earliest', # 启动策略
max_poll_records=1, # 每次拉取 1 条 → 等效 XREADGROUP count=1
max_poll_interval_ms=3600000, # 长任务不触发 rebalance (1h)
session_timeout_ms=120000, # 消费者组会话超时 (2m)
)
async def start(self):
await self._consumer.start()
async def poll(self, timeout_ms: int = 5000) -> QueueMessage | None:
records = await self._consumer.getmany(timeout_ms=timeout_ms, max_records=1)
for tp, msgs in records.items():
for m in msgs:
return QueueMessage(
message_id=f"{tp.partition}:{m.offset}",
payload=m.value,
topic=tp.topic,
partition=tp.partition,
offset=m.offset,
timestamp=m.timestamp,
_consumer_record=m, # 保存原始记录用于 commit
)
return None
async def commit(self, message: QueueMessage) -> bool:
"""commit offset = Redis XACK 等效操作"""
try:
await self._consumer.commit()
return True
except Exception as e:
logging.exception(f"Kafka commit failed: {e}")
return False
async def close(self):
await self._consumer.stop()
一个 Consumer Group (task_executor) 同时订阅所有任务 Topic。Kafka 原生支持一个 Consumer 订阅多个 Topic,Partition 级别的并行替代了 Redis 的双队列轮询:
| 场景 | Redis 实现 | Kafka 实现 |
|---|---|---|
| 优先级处理 | 两个队列分别轮询,优先 Priority 1 | 所有 Topic 统一订阅,parse.priority 分区数少 → 竞争少 → 更快消费 |
| 任务隔离 | 无隔离(同一 Stream) | 不同 Topic 物理隔离,dataflow 不影响 parse |
| 负载均衡 | Consumer Group 自动分配 | Partition 分配策略 (RoundRobin / CooperativeSticky) |
| 未确认重放 | 手动遍历 Pending | Offset 回退 / 不提交即自动重试 |
保留现有 Task 字典结构不变,仅增加 Kafka 元数据 header:
// Kafka Message Headers
{
"x-task-type": "parse", // 任务类型
"x-task-id": "uuid-xxxx", // 任务 ID
"x-tenant-id": "tenant-xxxx", // 租户 ID
"x-priority": "0", // 优先级
"x-source": "document_upload", // 来源
"x-timestamp": "1718123456789", // 生产时间戳
"x-retry-count": "0", // 重试次数
"x-schema-version": "1" // 消息格式版本
}
// Kafka Message Value (JSON)
{
"id": "task-uuid",
"doc_id": "doc-uuid",
"from_page": 0,
"to_page": 12,
"task_type": "parse",
"priority": 0,
"digest": "xxhash64-hex",
"progress": 0.0,
"begin_at": "2026-06-17 10:30:00",
"parser_id": "pdf",
// ... 其余现有字段
}
当前系统无 DLQ 机制。Kafka 迁移后引入标准的三层错误处理:
Message consumed
│
▼
┌─────────────┐
│ Attempt 1 │── success ──▶ commit() → DONE
└──────┬──────┘
│ fail
▼
┌─────────────┐
│ Attempt 2 │── success ──▶ commit() → DONE
└──────┬──────┘
│ fail
▼
┌─────────────┐
│ Attempt 3 │── success ──▶ commit() → DONE
└──────┬──────┘
│ fail (exceeds max_retries=3)
▼
┌──────────────┐
│ Send to DLQ │ → ragflow.task.dlq (保留现场,人工介入)
└──────────────┘
│ commit() original offset (避免阻塞)
handle_task() 函数中增加重试计数和 DLQ 逻辑。
DLQ 消息保留原始 metadata + stacktrace + 失败时间戳,支持后续批量重放。
| 指标 | 来源 | 告警阈值 |
|---|---|---|
| Consumer Lag | kafka.consumer:records-lag | > 1000 条持续 5 分钟 |
| Produce 错误率 | kafka.producer:record-error-rate | > 1% |
| Consume 错误率 | kafka.consumer:record-error-rate | > 1% |
| 消息处理延迟 (e2e) | 应用埋点 (produce_ts → consume_ts) | > 60s |
| DLQ 积压 | kafka.consumer:records-lag (DLQ topic) | > 10 条 |
| Partition 分布 | kafka.consumer:assigned-partitions | Rebalance 频率 > 0.1次/min |
| # | 文件 | 操作 | 改动说明 | 影响等级 |
|---|---|---|---|---|
| 1 | rag/utils/message_queue.py | + 新增 | 消息队列抽象层:QueueMessage, MessageQueueProducer, MessageQueueConsumer, MessageQueueFactory |
核心 |
| 2 | rag/utils/kafka_producer.py | + 新增 | 基于 aiokafka 的异步 Producer 实现,支持 acks=all、压缩、重试 | 核心 |
| 3 | rag/utils/kafka_consumer.py | + 新增 | 基于 aiokafka 的异步 Consumer 实现,手动 commit,支持 DLQ 转发 | 核心 |
| 4 | rag/utils/redis_conn.py | − 保留,标记废弃 | queue_product/queue_consumer/get_unacked_iterator/queue_info 等队列方法标记 deprecated;保留非队列方法 (KV、Lock、Sorted Set) |
中 |
| 5 | rag/svr/task_executor.py | ⚠ 改造 | (a) import: 新增 message_queue 导入 (b) collect(): Redis → Kafka Consumer (c) report_status(): 队列信息改为 Kafka lag 指标 (d) main(): 增加 Kafka Consumer 初始化/关闭 |
核心 |
| 6 | api/db/services/task_service.py | ⚠ 改造 | (a) queue_tasks(): Redis → MQ Producer (b) queue_dataflow(): Redis → MQ Producer |
核心 |
| 7 | api/db/services/document_service.py | ⚠ 改造 | queue_raptor_o_graphrag_tasks(): Redis → MQ Producer; get_queue_length(): Redis → Kafka admin |
核心 |
| 8 | api/db/joint_services/memory_message_service.py | ⚠ 改造 | Memory 任务生产: Redis → MQ Producer | 中 |
| 9 | common/constants.py | ⚠ 改造 | 新增 Kafka Topic 名称常量,保留原有 Redis 常量(兼容期) | 低 |
| 10 | common/settings.py | ⚠ 改造 | 新增 Kafka 配置加载函数: KAFKA_BOOTSTRAP_SERVERS, MQ_BACKEND, 保留 get_svr_queue_name() |
中 |
| 11 | pyproject.toml | ⚠ 改造 | 新增依赖: aiokafka>=0.8.0, kafka-python>=2.0.2 (admin CLI) |
中 |
| 12 | docker/docker-compose-base.yml | ⚠ 改造 | 新增 Kafka + Zookeeper (或 KRaft) 服务定义 | 中 |
| 13 | docker/docker-compose.yml | ⚠ 改造 | RAGFlow 服务依赖增加 Kafka 健康检查 | 低 |
| 14 | docker/.env | ⚠ 改造 | 新增环境变量: KAFKA_BOOTSTRAP_SERVERS, MQ_BACKEND, KAFKA_TOPIC_PREFIX |
低 |
| 15 | docker/service_conf.yaml.template | ⚠ 改造 | 新增 kafka 配置段: bootstrap_servers, security_protocol, sasl_mechanism | 低 |
| 16 | docker/entrypoint.sh | ⚠ 改造 | Kafka 就绪等待逻辑: 启动前检查 Kafka 连通性 | 低 |
# 新增 Kafka 服务 (KRaft 模式,无需 Zookeeper)
kafka:
image: bitnami/kafka:3.9
hostname: kafka
ports:
- "9092:9092" # Internal (PLAINTEXT)
- "${KAFKA_EXT_PORT:-9093}:9093" # External
environment:
# KRaft consensus
KAFKA_CFG_NODE_ID: 1
KAFKA_CFG_PROCESS_ROLES: broker,controller
KAFKA_CFG_CONTROLLER_QUORUM_VOTERS: 1@kafka:9094
KAFKA_CFG_CONTROLLER_LISTENER_NAMES: CONTROLLER
# Listeners
KAFKA_CFG_LISTENERS: PLAINTEXT://:9092,EXTERNAL://:9093,CONTROLLER://:9094
KAFKA_CFG_ADVERTISED_LISTENERS: PLAINTEXT://kafka:9092,EXTERNAL://${KAFKA_EXT_HOST:-localhost}:${KAFKA_EXT_PORT:-9093}
KAFKA_CFG_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT,EXTERNAL:PLAINTEXT,CONTROLLER:PLAINTEXT
KAFKA_CFG_INTER_BROKER_LISTENER_NAME: PLAINTEXT
# Topic defaults
KAFKA_CFG_AUTO_CREATE_TOPICS_ENABLE: "false" # 禁止自动创建 Topic
KAFKA_CFG_DEFAULT_REPLICATION_FACTOR: 1 # 单节点 = 1
KAFKA_CFG_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
KAFKA_CFG_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 1
# Retention
KAFKA_CFG_LOG_RETENTION_HOURS: 168 # 7 days default
volumes:
- kafka_data:/bitnami/kafka
networks:
- ragflow
restart: unless-stopped
healthcheck:
test: ["CMD", "kafka-topics.sh", "--bootstrap-server", "localhost:9092", "--list"]
interval: 15s
timeout: 10s
retries: 30
# 新建: docker/kafka_init.sh
#!/bin/bash
# Kafka Topic 初始化 — 在容器启动前执行
BOOTSTRAP="kafka:9092"
PARTITIONS=${KAFKA_DEFAULT_PARTITIONS:-4}
REPLICATION=${KAFKA_REPLICATION_FACTOR:-1}
declare -A TOPICS=(
["ragflow.task.parse"]="8"
["ragflow.task.parse.priority"]="4"
["ragflow.task.dataflow"]="4"
["ragflow.task.graphrag"]="2"
["ragflow.task.memory"]="2"
["ragflow.task.dlq"]="2"
)
echo "Waiting for Kafka..."
until kafka-topics.sh --bootstrap-server $BOOTSTRAP --list &>/dev/null; do
sleep 2
done
for topic in "${!TOPICS[@]}"; do
partitions=${TOPICS[$topic]}
if kafka-topics.sh --bootstrap-server $BOOTSTRAP --describe --topic "$topic" &>/dev/null; then
echo "Topic $topic already exists, skipping."
else
kafka-topics.sh --bootstrap-server $BOOTSTRAP \
--create --topic "$topic" \
--partitions "$partitions" \
--replication-factor "$REPLICATION"
echo "Created topic: $topic ($partitions partitions)"
fi
done
echo "Kafka topic initialization complete."
# 新增 Kafka 配置段
kafka:
bootstrap_servers: '${KAFKA_BOOTSTRAP_SERVERS:-kafka:9092}'
client_id: 'ragflow-server'
security_protocol: '${KAFKA_SECURITY_PROTOCOL:-PLAINTEXT}'
sasl_mechanism: '${KAFKA_SASL_MECHANISM:-PLAIN}'
sasl_username: '${KAFKA_SASL_USERNAME:-}'
sasl_password: '${KAFKA_SASL_PASSWORD:-}'
# 消息队列后端选择
mq:
backend: '${MQ_BACKEND:-redis}' # 'redis' 或 'kafka'
consumer_group: 'ragflow-task-executor' # Kafka Consumer Group ID
# ===== Kafka Configuration =====
MQ_BACKEND=kafka # redis | kafka
KAFKA_BOOTSTRAP_SERVERS=kafka:9092
KAFKA_EXT_PORT=9093
KAFKA_EXT_HOST=localhost
KAFKA_DEFAULT_PARTITIONS=4
KAFKA_REPLICATION_FACTOR=1 # 单节点=1, 生产=3
KAFKA_SECURITY_PROTOCOL=PLAINTEXT # PLAINTEXT | SASL_PLAINTEXT | SASL_SSL
# KAFKA_SASL_MECHANISM=PLAIN
# KAFKA_SASL_USERNAME=
# KAFKA_SASL_PASSWORD=
# 新增全局配置
MQ_BACKEND = os.getenv('MQ_BACKEND', 'redis')
KAFKA_BOOTSTRAP_SERVERS = os.getenv('KAFKA_BOOTSTRAP_SERVERS', 'localhost:9092')
KAFKA_CONSUMER_GROUP = os.getenv('KAFKA_CONSUMER_GROUP', 'ragflow-task-executor')
def get_mq_producer():
"""延迟初始化 MQ Producer"""
from rag.utils.message_queue import MessageQueueFactory
return MessageQueueFactory.get_producer()
def get_mq_consumer():
"""延迟初始化 MQ Consumer"""
from rag.utils.message_queue import MessageQueueFactory
return MessageQueueFactory.get_consumer(KAFKA_CONSUMER_GROUP)
# 保留兼容 (废弃标记)
def get_svr_queue_name(priority: int) -> str:
"""@deprecated: 仅在 MQ_BACKEND=redis 时使用"""
...
aiokafka>=0.8.0 依赖到 pyproject.tomlqueue_tasks() → 双写模式 (同时写 Redis 和 Kafka)queue_dataflow() → 双写模式queue_raptor_o_graphrag_tasks() → 双写模式task_executor.py 的 collect() 使用 Kafka Consumer (feature flag)report_status() 使用 Kafka lag 指标总计预估: 8 周 (2 人团队,含测试和灰度)
| # | 风险 | 等级 | 影响 | 缓解措施 |
|---|---|---|---|---|
| R1 | Kafka 学习曲线 团队需要掌握 Kafka 运维(Topic 管理、Partition 调优、Rebalance 排障) |
高 | 运维复杂度上升,故障排查时间增加 | KRaft 模式简化部署(无需 ZK);提供运维 Runbook;Phase 1 预留学习时间 |
| R2 | Consumer Rebalance 导致任务中断 长任务(GraphRAG 可运行数十分钟)可能触发 Rebalance 超时 |
高 | 正在执行的长任务被重新分配,可能导致重复消费 | 设置 max.poll.interval.ms=3600000 (1h);配置 session.timeout.ms=120000;使用 CooperativeSticky 分配策略 |
| R3 | 消息顺序性破坏 同一文档的多个页面分片任务需要有序消费 |
中 | 分片处理结果可能出现竞态条件 | 以 doc_id 作为 Partition Key,保证同一文档消息进入同一分区 |
| R4 | Kafka 服务不可用 Kafka Broker 宕机导致所有任务停滞 |
高 | 全系统任务处理中断 | 生产环境至少 3 Broker;配置 ISR min=2;Phase 4 保留 Redis 降级路径 |
| R5 | 消息重复消费 提交 Offset 失败 + 任务已完成 → 重复执行 |
中 | 文档被重新解析/分块,产生重复数据 | Consumer 端实现去重逻辑 (基于 task_id + digest 幂等检查);MySQL Task 表唯一约束 |
| R6 | Python 异步兼容性 aiokafka 与现有 asyncio 事件循环的兼容性 |
中 | 死锁或事件循环异常 | 充分集成测试;aiokafka 文档成熟,社区活跃;可 fallback 到 confluent-kafka (同步) |
| R7 | 部署复杂度增加 Docker Compose 增加 Kafka 服务 |
低 | 资源使用增加 ~1GB 内存 + 磁盘 | 单节点 KRaft 模式资源轻量;通过 profile 控制是否启动 |
| 当前 (Redis) | 迁移后 (Kafka) | |
|---|---|---|
| Python 客户端 | valkey==6.0.2 | valkey==6.0.2 (保留,KV/缓存)aiokafka>=0.8.0 (新增) |
| Docker 镜像 | valkey/valkey:8 | valkey/valkey:8 (保留)bitnami/kafka:3.9 (新增) |
| 依赖服务 | 1 (Redis) | 2 (Redis + Kafka) |
| 总内存预算 | ~128MB (Redis) | ~128MB (Redis) + ~1GB (Kafka heap) + ~256MB (OS page cache) |
| 磁盘预算 | ~1GB (Redis RDB/AOF) | ~1GB (Redis) + 20-50GB (Kafka log segments, 可调) |
| 文件 | 行数 | 角色 |
|---|---|---|
| rag/utils/redis_conn.py | 563 | Redis 连接层,Queue 生产者/消费者 API |
| rag/svr/task_executor.py | 1796 | Worker 主进程,任务收集和执行循环 |
| api/db/services/task_service.py | 555 | Task CRUD + queue_tasks/queue_dataflow 生产者 |
| api/db/services/document_service.py | 1109 | queue_raptor_o_graphrag_tasks 生产者 |
| common/constants.py | 259 | SVR_QUEUE_NAME, SVR_CONSUMER_GROUP_NAME |
| common/settings.py | ~140 | get_svr_queue_name(), Redis 配置加载 |
| docker/entrypoint.sh | ~340 | task_executor 启动逻辑 |
| docker/docker-compose-base.yml | ~241 | Redis 服务定义 |
| Redis Streams 概念 | Kafka 等价概念 |
|---|---|
| Stream (流) | Topic |
| Consumer Group | Consumer Group |
| Consumer | Consumer (per partition) |
| Message ID (时间戳-序号) | Offset (分区内递增整数) |
| Consumer Group PEL (Pending Entries List) | Uncommitted Offsets |
| XADD | Producer.send() |
| XREADGROUP | Consumer.poll() |
| XACK | Consumer.commit() |
| XPENDING | kafka-consumer-groups --describe |
| XGROUP CREATE | Group 自动创建 (首次 commit) |
| Stream TTL / MAXLEN | Topic Retention Policy (time/size) |
| N/A (无原生支持) | Partition (分区并行) |