RAGFlow 异步消息队列迁移至 Apache Kafka

技术分析报告 — 从 Redis Streams 到 Kafka 的完整改造方案

📅 版本: v0.25.6 📅 日期: 2026-06-17 📊 状态: 技术评估

1. 当前架构分析

1.1 整体架构概览

RAGFlow v0.25.6 使用 Redis (Valkey) Streams 作为异步任务消息代理。 这是一个基于 Redis Stream 消费者组的自研轻量级消息队列系统,无 Celery、RabbitMQ 或 Kafka 等第三方消息中间件。

┌──────────────────────────────────────────────────────────────────┐ │ RAGFlow 异步任务架构 │ │ │ │ ┌──────────┐ ┌──────────┐ ┌──────────┐ ┌──────────┐ │ │ │ HTTP │ │ Canvas │ │ Document │ │ Memory │ │ │ │ Upload │ │ Pipeline │ │ Service │ │ Service │ │ │ │ API │ │ API │ │ │ │ │ │ │ └────┬─────┘ └────┬─────┘ └────┬─────┘ └────┬─────┘ │ │ │ │ │ │ │ │ ▼ ▼ ▼ ▼ │ │ ┌────────────────────────────────────────────────────────┐ │ │ │ 📤 Producers │ │ │ │ queue_tasks() queue_dataflow() queue_raptor_...() │ │ │ │ [task_service] [task_service] [document_service] │ │ │ └────────────────────────┬───────────────────────────────┘ │ │ │ XADD │ │ ▼ │ │ ┌────────────────────────────────────────────────────────┐ │ │ │ Redis (Valkey) Streams │ │ │ │ ┌──────────────────────┐ ┌──────────────────────┐ │ │ │ │ │ rag_flow_svr_queue_1 │ │ rag_flow_svr_queue │ │ │ │ │ │ (Priority 1) │ │ (Priority 0) │ │ │ │ │ │ Consumer Group: │ │ Consumer Group: │ │ │ │ │ │ "rag_flow_svr_ │ │ "rag_flow_svr_ │ │ │ │ │ │ task_broker" │ │ task_broker" │ │ │ │ │ └──────────────────────┘ └──────────────────────┘ │ │ │ └────────────────────────┬───────────────────────────────┘ │ │ │ XREADGROUP │ │ ▼ │ │ ┌────────────────────────────────────────────────────────┐ │ │ │ 📥 Consumers (Task Executors) │ │ │ │ ┌──────────┐ ┌──────────┐ ┌──────────┐ │ │ │ │ │ Worker 0 │ │ Worker 1 │ │ Worker N │ │ │ │ │ │ max_conc │ │ max_conc │ │ max_conc │ │ │ │ │ │ =5 tasks │ │ =5 tasks │ │ =5 tasks │ │ │ │ │ └──────────┘ └──────────┘ └──────────┘ │ │ │ └────────────────────────────────────────────────────────┘ │ │ │ │ Task Types: parse/chunk | dataflow | graphrag | raptor | memory │ └──────────────────────────────────────────────────────────────────┘

1.2 核心组件详解

1.2.1 连接层 — rag/utils/redis_conn.py

单例 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重新入队消息

1.2.2 消费者 (Worker) — rag/svr/task_executor.py

核心异步 Worker 进程,关键参数:

1.2.3 生产者 (任务提交) — api/db/services/task_service.py

函数任务类型触发场景调用方
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

1.3 消息生命周期

1. Producer 构造 Task 字典 (含 id, doc_id, from_page, to_page, task_type, etc.) │ 2. ├─ 写入 MySQL (bulk_insert_into_db / replace_on_conflict) │ 3. ├─ 调用 REDIS_CONN.queue_product(queue_name, message=task) │ └─ Redis: XADD queue_name * message {...} │ 4. ▼ 消息进入 Redis Stream │ 5. Consumer (task_executor) 通过 XREADGROUP 拉取消息 │ 6. ├─ 从 MySQL 同步任务最新状态 (TaskService.get_task) │ 7. ├─ 检查取消标记 (has_canceled via Redis GET) │ 8. ├─ 执行具体任务逻辑 (build_chunks → embedding → insert_chunks) │ ├─ dataflow: run_dataflow() │ ├─ graphrag: run_graphrag_for_kb() │ ├─ raptor: run_raptor_for_kb() │ └─ memory: handle_save_to_memory_task() │ 9. ├─ 成功: redis_msg.ack() → XACK │ 10. └─ 失败: 消息保持 Pending 状态 (无 DLQ 机制)

1.4 当前架构成熟度评估

维度现状评级
消息持久化 Redis Stream 持久化 (AOF/RDB),但 128MB 内存上限触发 LRU 淘汰 ⚠️ 中等
高可用 单实例 Redis,无哨兵/集群 ❌ 低
水平扩展 Worker 可水平扩展,但所有 Worker 共享同一 Redis ⚠️ 中等
消息回溯 支持未确认消息重放,但不支持按时间/offset 回溯 ⚠️ 中等
死信队列 无 DLQ,失败消息永久 Pending ❌ 低
消息顺序 支持,但无分区概念 ✅ 良好
监控可观测 基础心跳 + Pending/Lag 计数,无消费延迟监控 ⚠️ 中等
消息吞吐 适合中等规模 (受限于单 Redis 实例) ⚠️ 中等
关键风险:Redis 配置 --maxmemory 128mb --maxmemory-policy allkeys-lru 意味着在高负载下,队列消息可能被 LRU 淘汰,导致任务丢失。生产环境不建议依赖此配置。

2. Kafka 迁移目标与收益

目标当前 (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
迁移核心收益:解决单点故障风险、消息持久化不足、无法水平扩展 Broker 三大痛点, 同时为未来事件溯源、实时数据管道、多租户隔离等高级特性奠定基础。

3. 改造技术方案

3.1 总体设计原则

  1. 渐进式迁移: 通过抽象层隔离 Kafka 细节,支持 Redis/Kafka 双模式切换
  2. 向后兼容: 保留现有 API 接口不变,仅替换底层传输
  3. 最小侵入: 对 Task 模型、API 层、业务逻辑层零改动
  4. 可观测优先: 所有消息生产/消费均埋点,支持 Prometheus 指标导出

3.2 Kafka Topic 规划

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 设计:key = f"{tenant_id}:{doc_id}:{task_id}"key = f"{tenant_id}:{flow_id}"(dataflow)。

3.3 Producer 层改造

3.3.1 新增消息队列抽象层

新建 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, ...)

3.3.2 Producer 调用点改造

所有现有 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)

3.3.3 Kafka Producer 实现

# 新建: 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()

3.4 Consumer 层改造

3.4.1 改造 collect() 函数

rag/svr/task_executor.pycollect() 函数从 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

3.4.2 Kafka Consumer 实现

# 新建: 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()

3.4.3 多 Topic 消费策略

一个 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)
未确认重放手动遍历 PendingOffset 回退 / 不提交即自动重试

3.5 消息格式设计

保留现有 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",
  // ... 其余现有字段
}

3.6 错误处理与死信队列

当前系统无 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 + 失败时间戳,支持后续批量重放。

3.7 监控与可观测性

指标来源告警阈值
Consumer Lagkafka.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-partitionsRebalance 频率 > 0.1次/min

4. 逐文件改造清单

#文件操作改动说明影响等级
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 连通性

5. Docker 部署变更

5.1 Kafka 服务定义 (docker-compose-base.yml)

# 新增 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

5.2 Topic 初始化脚本

# 新建: 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."

6. 配置变更

6.1 service_conf.yaml.template 新增配置段

# 新增 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

6.2 .env 新增环境变量

# ===== 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=

6.3 common/settings.py 新增

# 新增全局配置
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 时使用"""
    ...

7. 迁移实施路线图

Phase 1: 基础设施准备 (预计 2 周)

  • 新增 aiokafka>=0.8.0 依赖到 pyproject.toml
  • 实现 rag/utils/message_queue.py 抽象层
  • 实现 rag/utils/kafka_producer.pyrag/utils/kafka_consumer.py
  • Docker Compose 增加 Kafka 服务定义 (KRaft 模式,单节点)
  • Topic 初始化脚本 docker/kafka_init.sh
  • 配置模板和 .env 变量新增
  • 里程碑 M1: Kafka 基础组件可运行,单元测试通过

Phase 2: Producer 层改造 (预计 1 周)

  • 改造 queue_tasks() → 双写模式 (同时写 Redis 和 Kafka)
  • 改造 queue_dataflow() → 双写模式
  • 改造 queue_raptor_o_graphrag_tasks() → 双写模式
  • 改造 Memory 任务生产 → 双写模式
  • 里程碑 M2: 所有 Producer 点均支持 Kafka 写入,Redis 路径仍然有效

Phase 3: Consumer 层改造 (预计 2 周)

  • 改造 task_executor.pycollect() 使用 Kafka Consumer (feature flag)
  • 改造 report_status() 使用 Kafka lag 指标
  • 实现 DLQ 消费者逻辑
  • 增加 Kafka 消费者健康检查
  • 里程碑 M3: 集成测试通过,Kafka 消费者可正确处理所有任务类型

Phase 4: 灰度与切换 (预计 2 周)

  • 灰度部署:10% → 50% → 100% 流量切到 Kafka
  • 监控 Consumer Lag、吞吐量、错误率
  • 性能对比压测 (Redis vs Kafka 吞吐/延迟)
  • DR 演练: Kafka 故障时的降级回退方案
  • 里程碑 M4: 100% Kafka 流量,Redis 队列路径保留但关闭

Phase 5: 清理与优化 (预计 1 周)

  • 移除双写代码,统一 MQ_BACKEND=kafka
  • 清理 Redis 队列相关 dead code (标记 deprecated → 删除)
  • 更新文档、运维手册、监控面板
  • 里程碑 M5: 代码干净,文档齐全,系统稳定运行

时间线总览

Phase 1: 基础设施
Week 1 ──── Week 2 ──── Week 3 ──── Week 4 ──── Week 5 ──── Week 6 ──── Week 7 ──── Week 8
Phase 2: Producer
Phase 3: Consumer
Phase 4: 灰度切换
Phase 5: 清理优化

总计预估: 8 周 (2 人团队,含测试和灰度)

8. 风险评估与缓解

#风险等级影响缓解措施
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 控制是否启动
最高优先级规避项:Consumer Rebalance 导致的长任务中断 (R2) 和 Kafka 单点故障 (R4)。 这两个风险直接影响系统可用性,必须在 Phase 3 设计阶段充分验证。

9. 附录

9.1 当前系统消息流全景图

┌───────────────────────────┐ │ HTTP / API Layer │ │ ┌────────┐ ┌───────────┐ │ │ │Document│ │ Canvas │ │ │ │Upload │ │ Execute │ │ │ │API │ │ API │ │ │ └───┬────┘ └─────┬─────┘ │ │ │ │ │ └──────┼─────────────┼────────┘ │ │ ┌──────────────┼──────┐ ┌────┼──────────────┐ │ task_service.py │ │ task_service.py │ │ queue_tasks() │ │ queue_dataflow() │ │ (文档解析任务) │ │ (Canvas 工作流) │ └──────────┬──────────┘ └────┬──────────────┘ │ │ ┌──────────┼──────┐ ┌────────┼──────────────┐ │ document_service │ │ memory_message_ │ │ queue_raptor_o_ │ │ service.py │ │ graphrag_tasks() │ │ (Memory 任务) │ └──────────┬────────┘ └────────┬────────────┘ │ │ ▼ ▼ ┌──────────────────────────────────────────┐ │ REDIS_CONN (Singleton) │ │ queue_product(queue, msg) │ │ XADD │ └──────────────────┬───────────────────────┘ │ ┌────────┴────────┐ ▼ ▼ ┌──────────────────────┐ ┌──────────────────────┐ │ Redis Stream: │ │ Redis Stream: │ │ rag_flow_svr_queue_1 │ │ rag_flow_svr_queue │ │ (Priority 1) │ │ (Priority 0) │ └──────────┬───────────┘ └──────────┬───────────┘ │ │ │ XREADGROUP │ └────────────┬───────────┘ ▼ ┌──────────────────────────────────────────┐ │ task_executor.py (Worker) │ │ ┌────────────────────────────────────┐ │ │ │ collect() │ │ │ │ ├─ get_unacked_iterator() │ │ │ │ │ (重启恢复: 重放 Pending) │ │ │ │ └─ queue_consumer(pri_1, then_0) │ │ │ ├─ handle_task() │ │ │ │ └─ do_handle_task() │ │ │ │ ├─ parse → build_chunks() │ │ │ │ │ → embedding() │ │ │ │ │ → insert_chunks() │ │ │ │ ├─ dataflow → run_dataflow() │ │ │ │ ├─ graphrag → run_graphrag() │ │ │ │ ├─ raptor → run_raptor() │ │ │ │ └─ memory → handle_memory() │ │ │ └─ redis_msg.ack() (XACK) │ │ └────────────────────────────────────────┘ │ └──────────────────────────────────────────┘

9.2 依赖项对比

当前 (Redis)迁移后 (Kafka)
Python 客户端valkey==6.0.2valkey==6.0.2 (保留,KV/缓存)
aiokafka>=0.8.0 (新增)
Docker 镜像valkey/valkey:8valkey/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, 可调)

9.3 关键文件索引

文件行数角色
rag/utils/redis_conn.py563Redis 连接层,Queue 生产者/消费者 API
rag/svr/task_executor.py1796Worker 主进程,任务收集和执行循环
api/db/services/task_service.py555Task CRUD + queue_tasks/queue_dataflow 生产者
api/db/services/document_service.py1109queue_raptor_o_graphrag_tasks 生产者
common/constants.py259SVR_QUEUE_NAME, SVR_CONSUMER_GROUP_NAME
common/settings.py~140get_svr_queue_name(), Redis 配置加载
docker/entrypoint.sh~340task_executor 启动逻辑
docker/docker-compose-base.yml~241Redis 服务定义

9.4 术语对照表

Redis Streams 概念Kafka 等价概念
Stream (流)Topic
Consumer GroupConsumer Group
ConsumerConsumer (per partition)
Message ID (时间戳-序号)Offset (分区内递增整数)
Consumer Group PEL (Pending Entries List)Uncommitted Offsets
XADDProducer.send()
XREADGROUPConsumer.poll()
XACKConsumer.commit()
XPENDINGkafka-consumer-groups --describe
XGROUP CREATEGroup 自动创建 (首次 commit)
Stream TTL / MAXLENTopic Retention Policy (time/size)
N/A (无原生支持)Partition (分区并行)