Bella OpenAPI Video Generation Workflow: How Video Jobs Are Processed
Bella OpenAPI processes video generation as a two-stage asynchronous pipeline: jobs are first created and enqueued to Redis, then submitted to AI providers via scheduled executors, with continuous polling until finalization.
The lianjiatech/bella-openapi repository implements a robust, queue-driven architecture for handling video generation requests. This article breaks down the complete video generation workflow—from REST API ingestion through Redis-backed job distribution to provider-specific adaptor integration—using actual source file paths and implementation details from the codebase.
Stage 1: Job Creation and Enqueuing
REST API Entry Point
Client requests enter through VideoController.java at api/server/src/main/java/com/ke/bella/openapi/endpoints/VideoController.java. The controller exposes a POST /v1/video endpoint that accepts a VideoCreateRequest payload and delegates to the service layer.
Service Layer and Persistence
The VideoService.createVideoJob(...) method in api/server/src/main/java/com/ke/bella/openapi/service/VideoService.java orchestrates initial job setup:
- Generates a global
videoIdusingVideoIdGenerator.VIDEO_ID_GENERATOR.generate(spaceCode) - Persists a
VideoJobDBrow with statusqueued - Enqueues the ID into the Redis submit list
// Inside VideoService.createVideoJob
String videoId = VideoIdGenerator.VIDEO_ID_GENERATOR.generate(spaceCode);
videoJobDB.setStatus(Status.queued.name());
videoRepo.addVideoJob(videoJobDB);
queueManager.enqueueForSubmit(model, videoId); // Redis LPUSH
Redis Queue Integration
The VideoJobQueues class at api/server/src/main/java/com/ke/bella/openapi/queue/VideoJobQueues.java manages the submit queue (bella:video:submit:<model>) and sync queue (bella:video:syncing). Jobs remain in the submit queue until assigned to a provider channel.
Stage 2: Submit Queue Processing
The VideoJobExecutor Scheduler
VideoJobExecutor.java at api/server/src/main/java/com/ke/bella/openapi/executor/VideoJobExecutor.java starts after Spring Boot initialization (@PostConstruct). It runs two fixed-rate schedulers:
- Submit scheduler:
scheduleIntervalSeconds(default 5s) →processVideoJobs() - Sync scheduler:
syncIntervalSeconds(default 5s) →processSyncQueue()
RPM Throttling and Channel Selection
For each model, the executor:
- Acquires a distributed lock (
bella:video:model-lock:<model>) - Loads active channels via
ChannelService - Filters channels with available RPM quota using
ChannelRpmLimiteratapi/server/src/main/java/com/ke/bella/openapi/protocol/limiter/ChannelRpmLimiter.java - Calculates safe batch size via
calculateSafeBatchSizeByRpm
Batch Dequeue and Task Submission
The executor dequeues up to the calculated batch size from the submit queue:
List<String> videoIds = queueManager.dequeueForSubmit(model, batchSize);
submitBatchTasksToChannels(videoIds, availableChannels);
Each job is assigned to a channel (round-robin), the RPM quota is consumed, and a VideoJobSubmitTask is submitted to the TaskExecutor worker pool.
Stage 3: Provider Submission
VideoJobSubmitTask Execution
VideoJobSubmitTask.java at api/server/src/main/java/com/ke/bella/openapi/executor/VideoJobSubmitTask.java executes the following:
- Loads the job from
videoRepo.queryVideoJob - Invokes the adaptor via
VideoAdaptor.submitVideoTask(...) - Receives a channel-specific video ID from the provider
- Performs a CAS (compare-and-swap) update from
queued→submitting→processing - Enqueues the job ID to the sync queue (
VideoJobQueues.enqueueForSync)
Adaptor Pattern and Provider Integration
The VideoAdaptor interface at api/server/src/main/java/com/ke/bella/openapi/protocol/video/VideoAdaptor.java abstracts provider-specific implementations. The HuoshanAdaptor at api/server/src/main/java/com/ke/bella/openapi/protocol/video/HuoshanAdaptor.java converts OpenAPI requests to Huoshan format using HuoshanVideoConverter and manages authentication.
String channelVideoId = ((VideoAdaptor<VideoProperty>) adaptor)
.submitVideoTask(request, baseUrl, (VideoProperty) property, job.getVideoId());
Stage 4: Synchronization and Finalization
Sync Queue Polling
The sync scheduler in VideoJobExecutor.processSyncQueue() dequeues IDs from bella:video:syncing and submits VideoJobSyncTask instances to the worker pool.
VideoJobSyncTask and Provider Polling
VideoJobSyncTask.java at api/server/src/main/java/com/ke/bella/openapi/executor/VideoJobSyncTask.java handles the polling logic:
- Loads the job and verifies status is
processing - Checks
minProcessingSecondsthreshold; if not elapsed, re-enqueues - Queries provider status via
VideoAdaptor.queryVideoTask(...) - For terminal states (
completed,failed,cancelled):- Completed: Downloads file via
VideoAdaptor.transferVideoToFile(...)to OpenAI file service - Updates DB with size, duration, file ID, error JSON, progress
- Logs usage and cost via
EndpointLogger
- Completed: Downloads file via
ChannelVideoResult result = queryChannelVideoStatus(job, channel);
if (isTerminalState(result.getStatus())) {
handleTerminalState(result, channel); // download + processSyncResult
} else {
queueManager.enqueueForSync(videoId); // poll later
}
Cost Logging and Usage Tracking
VideoService.processSyncResult(...) in api/server/src/main/java/com/ke/bella/openapi/service/VideoService.java finalizes the job and invokes logVideoCost(...) to record provider usage data into EndpointProcessData, which is then forwarded to EndpointLogger at api/server/src/main/java/com/ke/bella/openapi/protocol/log/EndpointLogger.java for centralized billing analytics.
Job Lifecycle Management
Status States and Transitions
The VideoJobDB entity tracks jobs through the following states:
queued: Initial state after creation, waiting in submit queuesubmitting: Actively being sent to providerprocessing: Provider accepted job, generation in progresscompleted: Video generated successfully, file storedfailed: Generation error occurredcancelled: Job aborteddeleted: Soft-deleted by user
Transitions use CAS (compare-and-swap) logic to prevent race conditions during concurrent executor threads.
Deletion Constraints
Only jobs in terminal states (queued, completed, failed, cancelled) can be deleted. The VideoService.deleteVideoJob(...) method performs a soft delete by updating the status to deleted rather than removing the database row.
Code Examples
Creating a Video Job via cURL
curl -X POST https://api.example.com/v1/video \
-H "Authorization: Bearer <access_token>" \
-H "Content-Type: application/json" \
-d '{
"model":"huoshan-video-1.0",
"prompt":"A futuristic city skyline at sunrise",
"size":"720p",
"seconds":"10"
}'
The controller maps this JSON to VideoCreateRequest, then VideoService persists and enqueues the job.
Using the Java SDK
BellaOpenApiClient client = BellaOpenApiClient.builder()
.baseUrl("https://api.example.com")
.accessToken("YOUR_TOKEN")
.build();
VideoCreateRequest request = VideoCreateRequest.builder()
.model("huoshan-video-1.0")
.prompt("A kitten playing with a ball of yarn")
.size("720p")
.seconds("5")
.build();
VideoJobResponse response = client.video().create(request);
System.out.println("Video job created, id = " + response.getVideoId());
Behind the scenes, the SDK sends the same request that ends up in VideoService.createVideoJob.
Querying Job Status
curl -X GET https://api.example.com/v1/video/{videoId} \
-H "Authorization: Bearer <access_token>"
This endpoint reads VideoJobDB via VideoService.queryVideoJob and returns the current status.
Deleting a Finished Job
curl -X DELETE https://api.example.com/v1/video/{videoId} \
-H "Authorization: Bearer <access_token>"
This calls VideoService.deleteVideoJob, which updates the DB status to deleted.
Key Implementation Files
Summary
- Bella OpenAPI implements an asynchronous two-stage pipeline for video generation, separating job submission from result synchronization.
- Redis queues (
bella:video:submit:<model>andbella:video:syncing) decouple the REST API from provider communication, enabling horizontal scaling. - VideoJobExecutor manages both stages via scheduled threads, enforcing RPM limits through
ChannelRpmLimiterand distributed locking. - Adaptor pattern (
VideoAdaptor/HuoshanAdaptor) abstracts provider-specific protocols, allowing new video providers without changing core orchestration logic. - CAS state transitions ensure thread-safe updates from
queued→processing→completed/failed, with final cost logging viaEndpointLogger.
Frequently Asked Questions
How does Bella OpenAPI handle high concurrency for video generation?
Bella OpenAPI uses distributed Redis queues and per-model locking (bella:video:model-lock:<model>) to prevent race conditions. The VideoJobExecutor calculates safe batch sizes based on remaining RPM (requests per minute) quotas via ChannelRpmLimiter, ensuring channels never exceed provider rate limits. Jobs exceeding capacity remain in the Redis queue until capacity frees up.
What is the purpose of the VideoAdaptor interface?
The VideoAdaptor interface (defined in api/server/src/main/java/com/ke/bella/openapi/protocol/video/VideoAdaptor.java) abstracts provider-specific video APIs. Implementations like HuoshanAdaptor handle protocol conversion, authentication, and error mapping without modifying the core job orchestration. This pattern allows Bella OpenAPI to support multiple video providers (e.g., Huoshan, OpenAI, or custom endpoints) by simply adding new adaptor implementations.
How does the system ensure video jobs don't get lost during processing?
Jobs persist in PostgreSQL via VideoJobDB with strict status state machines (queued → submitting → processing → terminal). The VideoJobSubmitTask uses CAS (compare-and-swap) updates to ensure only valid state transitions occur. If a worker crashes, jobs remain in Redis queues (bella:video:submit:<model> or bella:video:syncing) until the VideoJobExecutor recovers them on the next scheduling interval (default 5 seconds).
Can I delete a video job while it's processing?
No. Bella OpenAPI restricts deletion to terminal states only: queued, completed, failed, or cancelled. The VideoService.deleteVideoJob(...) method performs a soft delete by updating the status to deleted rather than removing the record. This preserves audit trails and prevents accidental termination of active provider tasks that might incur costs.
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 →