How to Add Background Jobs and Async Processing in Go Clean Architecture
You can add background jobs to a Go Clean Architecture project by defining a JobDispatcher interface in the use-case layer, implementing an asynchronous worker pool in the adapter layer, and dispatching fire-and-forget tasks from within your use-case transactions without blocking the HTTP response.
The manakuro/golang-clean-architecture repository demonstrates strict separation of concerns across Domain, Use-case, Interface Adapter, and Infrastructure layers. When handling user creation in pkg/usecase/usecase/user.go, the code currently executes mailing, logging, and other processes synchronously inside a database transaction. Adding background jobs or async processing in Go allows these heavy or non-critical operations to run outside the request-response cycle while maintaining Clean Architecture boundaries.
Architectural Overview
The project follows Robert C. Martin’s Clean Architecture with four distinct layers:
- Domain:
pkg/domain/modelcontains core entities likeUser - Use-case:
pkg/usecase/usecaseholds application logic (e.g.,userUsecaseinpkg/usecase/usecase/user.go) - Interface Adapters:
pkg/adaptertranslates between frameworks and use-cases - Framework/Drivers:
pkg/infrastructurehandles routing, DB drivers, and external services
The repository provides a transactional DB abstraction via repository.DBRepository (defined in pkg/usecase/repository/db.go and implemented in pkg/adapter/repository/db.go). In the current Create method, the Transaction function wraps user persistence with commented placeholders for mailing and logging—the exact locations where asynchronous processing should be injected.
Where to Hook Asynchronous Processing
You have two strategic insertion points for background jobs that preserve the Dependency Rule (inner layers depend only on abstractions):
- Inside the transaction: Dispatch a job description after persisting the user but before committing. This ensures the job is only queued if the transaction succeeds, though the actual execution happens later.
- After the transaction: Enqueue the job once the transaction returns successfully, suitable for operations that do not require the transaction’s atomicity.
Both approaches keep the use-case agnostic of concrete queue implementations (channels, Redis, or cloud queues), satisfying Clean Architecture constraints.
Step 1: Define the Job Dispatcher Interface
Create a repository interface in the use-case layer to abstract job dispatching:
// pkg/usecase/repository/job.go
package repository
type JobDispatcher interface {
Dispatch(jobName string, payload interface{}) error
}
This interface resides in the inner layer, allowing business logic to request background work without knowing how it is executed.
Step 2: Implement an In-Process Worker Pool
In the adapter layer, implement the interface using a buffered channel and a fixed worker pool. This provides immediate asynchronous capability without external infrastructure dependencies:
// pkg/adapter/repository/job_queue.go
package repository
import (
"encoding/json"
"log"
)
type job struct {
name string
payload []byte
}
type asyncJobDispatcher struct {
queue chan job
workers int
}
func NewAsyncJobDispatcher(workers, queueSize int) JobDispatcher {
d := &asyncJobDispatcher{
queue: make(chan job, queueSize),
workers: workers,
}
d.start()
return d
}
func (d *asyncJobDispatcher) Dispatch(jobName string, payload interface{}) error {
data, err := json.Marshal(payload)
if err != nil {
return err
}
d.queue <- job{name: jobName, payload: data}
return nil
}
func (d *asyncJobDispatcher) start() {
for i := 0; i < d.workers; i++ {
go func(id int) {
for j := range d.queue {
log.Printf("[worker %d] handling %s", id, j.name)
// Real implementation switches on j.name to call appropriate handlers
}
}(i)
}
}
The NewAsyncJobDispatcher constructor initializes the worker pool immediately, ensuring jobs are consumed asynchronously as soon as the application starts.
Step 3: Inject the Dispatcher into the Use-Case
Modify the userUsecase struct in pkg/usecase/usecase/user.go to accept the dispatcher, then dispatch jobs within the transaction:
// pkg/usecase/usecase/user.go
type userUsecase struct {
userRepository repository.UserRepository
dBRepository repository.DBRepository
jobDispatcher repository.JobDispatcher // new field
}
func NewUserUsecase(r repository.UserRepository, d repository.DBRepository,
j repository.JobDispatcher) User {
return &userUsecase{r, d, j}
}
func (uu *userUsecase) Create(u *model.User) (*model.User, error) {
data, err := uu.dBRepository.Transaction(func(i interface{}) (interface{}, error) {
u, err := uu.userRepository.Create(u)
if err != nil {
return nil, err
}
// Fire-and-forget background jobs
_ = uu.jobDispatcher.Dispatch("welcomeEmail", u)
_ = uu.jobDispatcher.Dispatch("auditLog", u)
return u, nil
})
// Unchanged return handling
return data.(*model.User), err
}
The use-case now triggers background work without waiting for completion, keeping the HTTP response latency low while ensuring the jobs are only dispatched if the user is successfully created.
Step 4: Wire Dependencies in the Composition Root
Instantiate the dispatcher and inject it when building the use-case in your application entry point:
// cmd/app/main.go
func main() {
db, _ := gorm.Open("mysql", "...") // connection string omitted
dbRepo := repository.NewDBRepository(db)
// 5 workers, queue capacity 100
jobDispatcher := repository.NewAsyncJobDispatcher(5, 100)
userUC := usecase.NewUserUsecase(
repository.NewUserRepository(db),
dbRepo,
jobDispatcher,
)
// Inject userUC into controllers/router and start HTTP server
}
This wiring occurs in the outermost Framework layer, satisfying the dependency direction rule.
Architectural Benefits of This Approach
This design preserves Clean Architecture principles while enabling async processing:
- Dependency Rule: The use-case depends only on the
JobDispatcherinterface. Concrete implementations (channels, Redis, SQS) live in the adapter layer and can be swapped without touching business logic. - Testability: Unit tests can inject a mock dispatcher that records dispatch calls without spawning goroutines, allowing verification that jobs are queued without executing them.
- Scalability: The
asyncJobDispatchercan be replaced with a production-grade queue (RabbitMQ, AWS SQS, or Redis) by implementing the same interface, requiring zero changes topkg/usecase/usecase/user.go.
Summary
- Define a
JobDispatcherinterface inpkg/usecase/repository/job.goto abstract async operations - Implement the interface in
pkg/adapter/repository/job_queue.gousing channels and worker goroutines for in-process async handling - Inject the dispatcher into
userUsecasevia constructor and callDispatchinside theTransactionblock for fire-and-forget tasks - Wire the concrete dispatcher in
cmd/app/main.goas part of the dependency injection container - This pattern maintains Clean Architecture boundaries, supports testing with mocks, and allows migration to external message queues without refactoring use-case code
Frequently Asked Questions
Should background jobs run inside or outside the database transaction?
Dispatch jobs inside the transaction when you need to ensure the job is only queued if the database operation succeeds. In pkg/usecase/usecase/user.go, calling Dispatch before the transaction returns guarantees atomicity: if the user creation fails, no jobs are queued. For non-critical analytics or logging that should persist even if the transaction rolls back, dispatch after the transaction completes.
How do I test use-cases that dispatch background jobs?
Provide a mock implementation of repository.JobDispatcher in unit tests that records dispatch calls to a slice without spawning goroutines. Because the use-case depends only on the interface defined in pkg/usecase/repository/job.go, you can verify that Dispatch was called with the correct jobName and payload without executing real asynchronous logic or dealing with race conditions in tests.
Can I replace the in-process worker pool with Redis or RabbitMQ later?
Yes. The JobDispatcher interface in the inner use-case layer allows you to swap the asyncJobDispatcher implementation in pkg/adapter/repository/job_queue.go with a Redis-backed queue or cloud service (AWS SQS, GCP Pub/Sub) without modifying pkg/usecase/usecase/user.go. Simply create a new adapter that implements Dispatch to push to your external queue, and update the wiring in cmd/app/main.go.
What happens if a background job fails in the async worker pool?
The provided implementation uses fire-and-forget semantics (_ = uu.jobDispatcher.Dispatch(...)), meaning failures in json.Marshal return an error immediately, but runtime panics or processing errors in the worker goroutine must be handled within the worker’s execution logic. For production systems, implement retry logic, dead-letter queues, or structured logging inside the worker’s job handling switch statement in pkg/adapter/repository/job_queue.go.
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 →