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

> Discover the robust queue and DLQ design in Tencent WeKnora that scales Wiki ingest to 40k documents per knowledge base using Redis and Asynq for efficient processing and error handling.

- Repository: [Tencent/WeKnora](https://github.com/tencent/WeKnora)
- Tags: architecture
- Published: 2026-09-12

---

**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`](https://github.com/Tencent/WeKnora/blob/main/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`](https://github.com/Tencent/WeKnora/blob/main/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`](https://github.com/Tencent/WeKnora/blob/main/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`](https://github.com/Tencent/WeKnora/blob/main/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`](https://github.com/Tencent/WeKnora/blob/main/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:

```yaml

# 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

```go
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

```go
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`](https://github.com/Tencent/WeKnora/blob/main/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`](https://github.com/Tencent/WeKnora/blob/main/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.