Queue and DLQ Design for Scaling Wiki Ingest to 40k Documents per Knowledge Base in WeKnora

WeKnora scales the Wiki-ingest pipeline to approximately 40,000 documents per knowledge base by leveraging a dedicated Redis-backed Asynq task queue paired with a service-level dead-letter queue (DLQ) that archives exhausted tasks for manual triage.

The Tencent/WeKnora repository implements a high-throughput ingestion system designed to process tens of thousands of Wikipedia documents per knowledge base without interfering with chat or post-processing workloads. This examination of the queue and DLQ design for scaling Wiki ingest reveals how the architecture uses isolated worker pools, persistent task storage, and per-KB concurrency limits to maintain stability at scale.

Redis-Backed Task Queue Architecture

At the core of the ingestion pipeline is a dedicated Wiki queue implemented using the Asynq task queue library with Redis as the primary broker. The system defines specific constants in internal/types/task.go to isolate high-volume ingest operations:

  • QueueWiki – The named queue (wiki:ingest) routed to a dedicated worker pool
  • WorkerPoolWiki – An isolated pool with a default concurrency of 8 workers (DefaultWikiWorkerConcurrency)
  • TypeWikiIngest – The task type identifier (wiki:ingest) used for routing and metrics

This isolation prevents CPU-intensive document processing from starving interactive chat requests or background post-process tasks. According to the source code in internal/types/task.go, the worker pool guarantees capacity reservation through dedicated goroutine workers rather than shared thread pools.

Persistent Task Storage and Durability

While Redis provides the hot queue for task distribution, durability is guaranteed through a persistent PostgreSQL mirror described in the architecture documentation. The table task_pending_ops functions as a "Redis list queue persistent replacement," ensuring that tasks survive broker restarts or failovers.

Batch Claiming with Deduplication

To prevent a single document from being processed by two concurrent batches, the implementation uses a batch claim mechanism with row-level locking. The ClaimBatch operation groups operations by a dedup_key (the document ID) and employs FOR UPDATE SKIP LOCKED to atomically lock rows without blocking concurrent consumers. This pattern, documented in website-docs/02-architecture/05-async-tasks.md, ensures that exactly one worker claims a batch of operations for a specific document while other workers skip to available work.

Concurrency Control Per Knowledge Base

The system implements per-KB throttling to prevent any single knowledge base from monopolizing cluster resources. Configuration in internal/types/wiki_page.go exposes four key parameters through the WikiPageConfig struct:

  • IngestBatchSize – Number of operations processed per claim cycle
  • IngestMapParallel – Parallelism for the map phase of document processing
  • IngestReduceParallel – Parallelism for the reduce/aggregation phase
  • IngestMaxInflight – Maximum concurrent batches allowed for a specific KB

These limits ensure that a 40,000-document ingestion job proceeds at a predictable rate without overwhelming downstream vector databases or embedding services.

Dead-Letter Queue Implementation

When tasks exhaust their retry budget, they are archived to a service-level DLQ rather than dropped. The internal/types/task_dead_letter.go file defines the TaskDeadLetter struct, which persists:

  • TaskType and Scope (e.g., TaskScopeKnowledgeBase)
  • ScopeID (the KB identifier) and RelatedID (the specific document)
  • Payload and LastError for debugging
  • FailCount tracking the number of retry attempts

Two pathways feed the DLQ:

  1. Asynq's built-in middleware – Handles generic task failures with exponential backoff (approximately 10 seconds to 2.5 minutes)
  2. Service-level retry handling – When a Wiki-ingest operation exceeds wikiMaxFailRetries, the ingest service explicitly writes to the DLQ table

Finalization Locking

To guarantee exactly-once finalization per knowledge base, the system uses a Redis SETNX lock with the key pattern wiki:active:<kbID>. This prevents race conditions where multiple workers might attempt to mark a large ingestion job as complete simultaneously.

Configuration Example

The following YAML configuration demonstrates tuning parameters for a high-volume knowledge base:


# wiki_config.yaml for 40k document KB

auto_ingest: true
synthesis_model_id: "anthropic/claude-3"
max_pages_per_ingest: 0      # unlimited pages

ingest_batch_size: 5
ingest_map_parallel: 10
ingest_reduce_parallel: 10
ingest_max_inflight: 4      # limit concurrent batches

Code Examples

Enqueueing a Wiki Ingest Task

import (
    "encoding/json"
    "time"
    "github.com/hibiken/asynq"
    "github.com/Tencent/WeKnora/internal/types"
)

func EnqueueWikiIngest(client *asynq.Client, kbID string) error {
    payload, _ := json.Marshal(map[string]string{"kb_id": kbID})
    task := asynq.NewTask(types.TypeWikiIngest, payload)
    
    _, err := client.Enqueue(task,
        asynq.Queue(types.QueueWiki),
        asynq.MaxRetry(10),
        asynq.Timeout(30*time.Second),
    )
    return err
}

Querying Dead-Letter Entries

import (
    "github.com/Tencent/WeKnora/internal/types"
    "gorm.io/gorm"
)

func ListWikiDLQ(db *gorm.DB, kbID string) ([]types.TaskDeadLetter, error) {
    var letters []types.TaskDeadLetter
    err := db.
        Where("task_type = ?", types.TypeWikiIngest).
        Where("scope = ?", types.TaskScopeKnowledgeBase).
        Where("scope_id = ?", kbID).
        Find(&letters).Error
    return letters, err
}

Summary

  • Dedicated Worker Pools – The WorkerPoolWiki with configurable concurrency isolates ingestion workloads from interactive services.
  • PostgreSQL Durability – The task_pending_ops table mirrors the Redis queue, ensuring tasks survive broker restarts.
  • Deduplication – ClaimBatch uses dedup_key and FOR UPDATE SKIP LOCKED to prevent duplicate document processing.
  • Per-KB Limits – IngestMaxInflight and related configuration parameters prevent resource starvation across knowledge bases.
  • Dual-Path DLQ – Failed tasks route through either Asynq middleware or explicit service logic to TaskDeadLetter for observability and manual replay.

Frequently Asked Questions

How does WeKnora prevent the same Wiki document from being processed twice?

The ClaimBatch function in the ingestion service groups pending operations by a dedup_key corresponding to the document ID. When claiming work, it executes a SELECT ... FOR UPDATE SKIP LOCKED query against task_pending_ops, which atomically locks the rows for that specific document while allowing other workers to skip to unlocked rows. This guarantees that only one batch processor handles a given document at any time.

What happens when a Wiki ingest task fails permanently after retries?

Tasks that exceed the maximum retry count (configured via wikiMaxFailRetries) are written to the TaskDeadLetter table defined in internal/types/task_dead_letter.go. The system captures the LastError, FailCount, and full Payload, enabling operators to inspect failures via SQL queries and manually re-enqueue tasks after fixing underlying issues such as rate limits or invalid document formats.

How is concurrency controlled for knowledge bases with 40,000 documents?

The WikiPageConfig struct in internal/types/wiki_page.go exposes IngestMaxInflight, which limits the number of concurrent batches active for a specific KB. Combined with IngestMapParallel and IngestReduceParallel, these settings throttle the pipeline to maintain stable resource usage regardless of total document count, preventing downstream services from being overwhelmed during large ingestions.

Why does WeKnora use both Redis and PostgreSQL for the task queue?

Redis provides the low-latency, high-throughput broker necessary for Asynq's distributed worker coordination, while PostgreSQL's task_pending_ops table acts as a persistent ledger. This hybrid approach ensures that if the Redis instance restarts or loses data, the pending operations can be recovered from the relational database, maintaining the durability guarantees required for production knowledge base ingestion at scale.

Have a question about this repo?

These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:

Share the following with your agent to get started:
curl -s "https://instagit.com/install.md"

Works with
Claude Codex Cursor VS Code OpenClaw Any MCP Client

Maintain an open-source project? Get it listed too →