Understanding the Job Queue System in Bella OpenAPI: How `@EnableJobQueue` Works

Bella OpenAPI implements a Redis-backed distributed job queue system that decouples long-running tasks from synchronous request handling, activated by the @EnableJobQueue annotation which imports JobQueueConfiguration to register JobQueueProperties and enable the queue subsystem.

The Bella OpenAPI repository (lianjiatech/bella-openapi) provides a scalable infrastructure for AI service orchestration. Its job queue system allows controllers to enqueue tasks such as video processing or audio transcription, while background workers dequeue and execute them asynchronously using Redis as the persistent storage layer.

Architecture of the Bella OpenAPI Job Queue System

The system follows a layered architecture that separates configuration, abstraction, storage implementation, and domain-specific usage.

Configuration Layer: @EnableJobQueue and JobQueueConfiguration

The entry point for the job queue system is the @EnableJobQueue annotation defined in api/spi/src/main/java/com/ke/bella/job/queue/config/EnableJobQueue.java. This annotation imports JobQueueConfiguration, which registers the configuration bean:

@Target(ElementType.TYPE)
@Retention(RetentionPolicy.RUNTIME)
@Documented
@Import({ JobQueueConfiguration.class })
public @interface EnableJobQueue { }

JobQueueConfiguration in api/spi/src/main/java/com/ke/bella/job/queue/config/JobQueueConfiguration.java declares a bean bound to the bella.job-queue namespace:

@Configuration
public class JobQueueConfiguration {

    @Bean
    @ConfigurationProperties(value = "bella.job-queue")
    public JobQueueProperties jobQueueProperties() {
        return new JobQueueProperties();
    }
}

The properties class in api/spi/src/main/java/com/ke/bella/job/queue/config/JobQueueProperties.java holds the remote service URL and timeout settings:

@Data
public class JobQueueProperties {
    private String url;                // e.g. http://job-queue-service
    private Integer defaultTimeout = 300;
}

Core Abstraction: The JobQueue<T> Interface

The generic contract for queue operations is defined in api/server/src/main/java/com/ke/bella/openapi/queue/JobQueue.java. It provides three primitive operations that map to Redis deque semantics (HEAD = LEFT, TAIL = RIGHT):

public interface JobQueue<T> {
    void lpush(String queueKey, T jobId);          // priority / retry (push to HEAD)
    void rpush(String queueKey, T jobId);          // new tasks (push to TAIL)
    List<T> lpop(String queueKey, int maxSize);    // dequeue batch from HEAD
}

Redis Implementation: RedisJobQueue<T>

The concrete implementation in api/server/src/main/java/com/ke/bella/openapi/queue/RedisJobQueue.java uses Redisson's RDeque to provide atomic operations across distributed JVM instances:

@Component
@Slf4j
public class RedisJobQueue<T> implements JobQueue<T> {

    @Resource
    private RedissonClient redissonClient;

    @Override
    public void lpush(String queueKey, T jobId) {
        RDeque<T> queue = redissonClient.getDeque(queueKey);
        queue.addFirst(jobId);
        log.debug("[JobQueue] lpush {} to HEAD: queue={}", jobId, queueKey);
    }

    @Override
    public void rpush(String queueKey, T jobId) {
        RDeque<T> queue = redissonClient.getDeque(queueKey);
        queue.addLast(jobId);
        log.debug("[JobQueue] rpush {} to TAIL: queue={}", jobId, queueKey);
    }

    @Override
    public List<T> lpop(String queueKey, int maxSize) {
        RDeque<T> queue = redissonClient.getDeque(queueKey);
        List<T> result = new ArrayList<>();
        for (int i = 0; i < maxSize; i++) {
            T jobId = queue.pollFirst();
            if (jobId == null) break;
            result.add(jobId);
        }
        if (!result.isEmpty())
            log.debug("[JobQueue] lpop {} jobs from HEAD: queue={}", result.size(), queueKey);
        return result;
    }
}

Domain-Specific Facade: VideoJobQueues

Business components wrap the generic queue with domain-specific naming rules. The VideoJobQueues class in api/server/src/main/java/com/ke/bella/openapi/queue/VideoJobQueues.java demonstrates this pattern:

@Component
public class VideoJobQueues {

    @Resource
    private JobQueue<String> jobQueue;

    private static final String SUBMIT_QUEUE_PREFIX = "bella:video:submit:";
    private static final String SYNCING_QUEUE_KEY   = "bella:video:syncing";

    public void enqueueForSubmit(String model, String videoId) {
        jobQueue.rpush(SUBMIT_QUEUE_PREFIX + model, videoId);
    }

    public void enqueueForSubmitFirst(String model, String videoId) {
        jobQueue.lpush(SUBMIT_QUEUE_PREFIX + model, videoId);
    }

    public List<String> dequeueForSubmit(String model, int maxSize) {
        return jobQueue.lpop(SUBMIT_QUEUE_PREFIX + model, maxSize);
    }

    public void enqueueForSync(String videoId) {
        jobQueue.rpush(SYNCING_QUEUE_KEY, videoId);
    }

    public List<String> dequeueForSync(int maxSize) {
        return jobQueue.lpop(SYNCING_QUEUE_KEY, maxSize);
    }
}

How @EnableJobQueue Activates the System

The activation mechanism follows Spring Boot's annotation-driven configuration pattern. When you add @EnableJobQueue to your main application class, as seen in api/server/src/main/java/com/ke/bella/openapi/Application.java, the following sequence occurs:

  1. Annotation Detection: Spring detects @EnableJobQueue on the application class
  2. Configuration Import: The @Import({ JobQueueConfiguration.class }) meta-annotation loads JobQueueConfiguration
  3. Bean Registration: JobQueueConfiguration registers JobQueueProperties as a bean bound to bella.job-queue configuration properties
  4. Component Scanning: Spring discovers RedisJobQueue (annotated with @Component) and injects the RedissonClient
@SpringBootApplication
@EnableMethodCache(basePackages = "com.ke.bella.openapi")
@EnableScheduling
@EnableJobQueue               // Activates the queue subsystem
public class Application { 
    // Application entry point
}

Using the Job Queue in Practice

Enqueuing Jobs

Services enqueue tasks using the domain-specific façade. For standard processing, use rpush (add to tail); for priority or retry scenarios, use lpush (add to head):

@Autowired
private VideoJobQueues videoJobQueues;

public void submitVideo(String model, String videoId) {
    // Standard enqueue - adds to tail (FIFO)
    videoJobQueues.enqueueForSubmit(model, videoId);
}

public void retryVideo(String model, String videoId) {
    // Priority enqueue - adds to head (LIFO for retries)
    videoJobQueues.enqueueForSubmitFirst(model, videoId);
}

Dequeuing Jobs

Workers typically run as scheduled tasks or background threads. They dequeue batches from the head of the queue using lpop to maximize throughput:

@Scheduled(fixedDelayString = "${bella.job-queue.poll-interval:5000}")
public void processVideoJobs() {
    // Dequeue up to 10 jobs at once
    List<String> batch = videoJobQueues.dequeueForSubmit("model-name", 10);
    
    for (String videoId : batch) {
        // Process each video job
        processVideo(videoId);
    }
}

Direct Queue Interface Usage

For advanced scenarios, inject the generic JobQueue<T> directly:

@Resource
private JobQueue<String> jobQueue;

public void customEnqueue(String queueName, String jobId) {
    jobQueue.rpush(queueName, jobId);
}

public List<String> customDequeue(String queueName, int batchSize) {
    return jobQueue.lpop(queueName, batchSize);
}

Summary

  • Bella OpenAPI implements a Redis-backed distributed job queue using Redisson's atomic deque operations to ensure consistency across multiple JVM instances.
  • The @EnableJobQueue annotation activates the subsystem by importing JobQueueConfiguration, which registers JobQueueProperties and enables component scanning for queue implementations.
  • The JobQueue<T> interface provides three core operations: lpush (head/priority), rpush (tail/FIFO), and lpop (batch dequeue).
  • RedisJobQueue<T> implements this interface using RDeque from Redisson, providing thread-safe, distributed queue operations.
  • Domain-specific facades like VideoJobQueues wrap the generic interface with business logic and naming conventions, separating infrastructure from business code.

Frequently Asked Questions

What is the purpose of the @EnableJobQueue annotation in Bella OpenAPI?

The @EnableJobQueue annotation serves as the configuration entry point for the job queue subsystem. When added to a Spring Boot application class, it imports JobQueueConfiguration, which registers the JobQueueProperties bean and enables the infrastructure components (like RedisJobQueue) to be discovered by component scanning. Without this annotation, the queue beans and configuration properties remain inactive.

How does Bella OpenAPI ensure thread safety in distributed job queue operations?

Thread safety and atomicity are guaranteed through Redisson, a Redis client for Java. The RedisJobQueue implementation uses Redisson's RDeque interface, which provides atomic operations like addFirst, addLast, and pollFirst. These operations are executed as atomic Redis commands, ensuring that even when multiple JVM instances (workers) dequeue jobs simultaneously, each job is processed by exactly one consumer without race conditions.

What is the difference between lpush and rpush in the Bella OpenAPI job queue?

In the JobQueue<T> interface, rpush adds elements to the tail (right side) of the deque, implementing a standard FIFO (First-In-First-Out) queue for new tasks. lpush adds elements to the head (left side), implementing a LIFO (Last-In-First-Out) pattern typically used for priority tasks or retry logic. Workers always dequeue from the head using lpop, so lpush jobs are processed before rpush jobs.

Can I use the job queue system for custom business logic beyond video processing?

Yes. While Bella OpenAPI provides VideoJobQueues as a domain-specific example, you can create custom facades by implementing the same pattern. Inject the generic JobQueue<T> bean (which resolves to RedisJobQueue at runtime) and define your own queue key naming conventions and wrapper methods. This allows you to create dedicated queues for image generation, data export, or any other asynchronous processing needs while leveraging the same Redis-backed infrastructure.

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 →