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 poolWorkerPoolWiki– 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 cycleIngestMapParallel– Parallelism for the map phase of document processingIngestReduceParallel– Parallelism for the reduce/aggregation phaseIngestMaxInflight– 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:
TaskTypeandScope(e.g.,TaskScopeKnowledgeBase)ScopeID(the KB identifier) andRelatedID(the specific document)PayloadandLastErrorfor debuggingFailCounttracking the number of retry attempts
Two pathways feed the DLQ:
- Asynq's built-in middleware – Handles generic task failures with exponential backoff (approximately 10 seconds to 2.5 minutes)
- 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
WorkerPoolWikiwith configurable concurrency isolates ingestion workloads from interactive services. - PostgreSQL Durability – The
task_pending_opstable mirrors the Redis queue, ensuring tasks survive broker restarts. - Deduplication –
ClaimBatchusesdedup_keyandFOR UPDATE SKIP LOCKEDto prevent duplicate document processing. - Per-KB Limits –
IngestMaxInflightand related configuration parameters prevent resource starvation across knowledge bases. - Dual-Path DLQ – Failed tasks route through either Asynq middleware or explicit service logic to
TaskDeadLetterfor 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:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →