Complete guide for the Redis Queue (RQ) integration between Opik's Java backend and Python workers using the official RQ library.
Status: ✅ Working end-to-end (Plain JSON contract; no custom serializer)
Last Updated: 2025-10-15
- ✅ Java RqPublisher - Creates RQ-compatible Redis HASH structures
- ✅ Plain JSON
datafield - UTF-8 JSON (no compression) - ✅ RQ-native Redis structure - Keys and lists match RQ defaults (e.g.,
rq:queue:<queue>) - ✅ Python RQ Worker via RqWorkerManager - Starts under Gunicorn with JSONSerializer + default Job
- ✅ OpenTelemetry Metrics - Metrics emitted from
MetricsWorker - ✅ Robust Connection Management - Exponential backoff retry logic
- ✅ Aligned Logging - Unified format with pid/process and thread info
- Switched from zlib-compressed
datato plain JSON (UTF-8) - Removed custom serializer/job; using RQ's
JSONSerializerand defaultJob - Pre-consume "func injection" removed (RQ restores from
datapayload) - No-op death penalty used to avoid signals in background thread
- Queue key corrected to
rq:queue:<queue-name>
- ✅ Redis: localhost:6379 (password:
opik) - ✅ MySQL: localhost:3306
- ✅ ClickHouse: localhost:8123
cd apps/opik-python-backend
source venv/bin/activate
export REDIS_HOST=localhost REDIS_PORT=6379 REDIS_DB=0 REDIS_PASSWORD=opik
python src/opik_backend/rq_worker.pycd apps/opik-backend
java -jar target/opik-backend-1.0-SNAPSHOT.jar server config.yml# Send message
curl -X POST "http://localhost:8080/v1/internal/hello-world?message=Test"
# Check queue
curl http://localhost:8080/v1/internal/hello-world/queue-size- Overview
- Architecture
- Detailed Setup
- Components
- OpenTelemetry Metrics
- Queue Configuration
- Usage Guide
- Adding New Queues
- Testing
- Troubleshooting
- Design Decisions
- Refactoring History
This integration enables the Java backend to enqueue jobs that are processed asynchronously by Python workers using Redis Queue (RQ). Production path uses RQ-native contracts (plain JSON) without Python bridges or custom serializers. This is useful for:
- CPU-intensive Python tasks (ML inference, data processing)
- Python-specific libraries (optimizer, analytics)
- Async job processing (background tasks, scheduled jobs)
- Scaling independently (Java services and Python workers)
- ✅ Type-safe queue definitions using Java enums
- ✅ Immutable message format using Java records
- ✅ Configuration-driven TTL management
- ✅ Interface-based design for testability
- ✅ Full RQ protocol compatibility
- ✅ Multiple queue support
Java directly creates RQ-compatible job structures in Redis for processing by Python RQ workers. The data field is plain JSON (UTF-8). The worker uses RQ's default JSONSerializer and default Job.
┌─────────────────────────────────────────────────────────────────┐
│ Java Backend (Redisson) │
│ - Creates RQ-compatible Redis HASH │
│ - Stores: created_at, enqueued_at, status, origin, timeout │
│ - Stores: data (plain JSON [func, null, args, {}]) │
│ - Adds job ID to Redis list (queue) │
└──────────────────────────┬──────────────────────────────────────┘
│
▼
┌─────────────────────────────────────────────────────────────────┐
│ Redis Server │
│ - Job data: rq:job:{id} (Redis HASH, RQ format) │
│ - Queue list: rq:queue:opik:optimizer-cloud │
└──────────────────────────┬──────────────────────────────────────┘
│
▼
┌─────────────────────────────────────────────────────────────────┐
│ RQ Worker (Python) │
│ - Uses JSONSerializer (default) │
│ - Default Job class │
│ - Configured with decode_responses=False │
│ - ✅ Processes jobs end-to-end │
└─────────────────────────────────────────────────────────────────┘
What Works:
- ✅ Java creates RQ-compatible Redis HASH structures
- ✅ Plain JSON
dataarray[func, null, args, kwargs] - ✅ Redis structure identical to Python-created jobs
- ✅
Job.fetch()and RQ worker processing succeed with JSONSerializer - ✅ End-to-end processing via
RqWorkerManagerin production
┌─────────────────┐ ┌─────────┐ ┌──────────────────┐
│ Java Backend │────────▶│ Redis │◀────────│ Python Worker │
│ (Producer) │ │ Queue │ │ (Consumer) │
│ │ │ │ │ │
│ RqPublisher │ RPUSH │ List │ LPOP │ RQ Worker │
│ QueueProducer │────────▶│ +Bucket│◀────────│ process_xxx() │
└─────────────────┘ └─────────┘ └──────────────────┘
RQ uses a two-tier storage approach:
- Job Metadata: Stored in
rq:job:{job-id}as a hash with full job details - Queue List: Contains only job IDs in a Redis list for FIFO processing
Redis Storage:
┌──────────────────────────────────────┐
│ rq:job:123abc (Hash) │
│ ├─ func: "process_optimizer_job" │
│ ├─ args: ["data"] │
│ ├─ status: "queued" │
│ └─ enqueued_at: "2025-10-14..." │
└──────────────────────────────────────┘
┌──────────────────────────────────────┐
│ rq:queue:opik:optimizer-cloud (List) │
│ ├─ "123abc" │
│ ├─ "456def" │
│ └─ "789ghi" │
└──────────────────────────────────────┘
com.comet.opik.infrastructure
├── queues/ # Queue abstractions
│ ├── QueueProducer.java # Interface for queue producers
│ ├── Queue.java # Enum of available queues
│ ├── RqMessage.java # Immutable message record
│ ├── RqQueueConfig.java # Queue configuration
│ └── JobStatus.java # Job status enum
├── redis/ # Redis implementation
│ └── RqPublisher.java # RQ implementation of QueueProducer
└── QueuesConfig.java # Configuration class
- Java 21+
- Python 3.8+
- Redis 7.x
- Maven 3.x
# Using Docker
docker run -d -p 6379:6379 --name opik-redis redis:7.2-alpine
# Or use existing Docker Compose
cd deployment/docker-compose
docker-compose up -d redisEdit apps/opik-backend/config.yml:
queues:
enabled: true
defaultJobTtl: 1 day
queues:
opik:optimizer-cloud:
jobTTl: 1 daycd apps/opik-backend
mvn clean package -DskipTestscd apps/opik-python-backend
# Install dependencies
pip install -r requirements.txt
# Set environment variables
export REDIS_HOST=localhost
export REDIS_PORT=6379
export REDIS_DB=0
# Start worker
python src/opik_backend/rq_worker.pyExpected output:
2025-10-14 10:00:00 INFO [opik_backend.rq_worker] - Starting RQ worker...
2025-10-14 10:00:00 INFO [opik_backend.rq_worker] - Connecting to Redis at localhost:6379 (db=0)
2025-10-14 10:00:00 INFO [opik_backend.rq_worker] - Listening on queues: ['opik:hello_world_queue', 'opik:optimizer-cloud']
2025-10-14 10:00:00 INFO [opik_backend.rq_worker] - RQ Worker started successfully
cd apps/opik-backend
java -jar target/opik-backend-1.0-SNAPSHOT.jar server config.yml# Send a test message
curl -X POST "http://localhost:8080/v1/internal/hello-world?message=Hello%20from%20Java"
# Response:
{
"status": "success",
"message": "Message enqueued successfully",
"jobId": "a1b2c3d4-e5f6-7890-abcd-ef1234567890",
"queue": "opik:optimizer-cloud",
"sentMessage": "Hello from Java"
}
# Check queue size
curl http://localhost:8080/v1/internal/hello-world/queue-size
# Response:
{
"queue": "opik:optimizer-cloud",
"size": 0
}Check Python worker logs:
2025-10-14 10:01:00 INFO [opik_backend.rq_worker] - Processing optimizer job: Hello from Java
2025-10-14 10:01:00 INFO [opik_backend.rq_worker] - Optimizer job processed successfully: {...}
Location: com.comet.opik.infrastructure.queues.QueueProducer
public interface QueueProducer {
/**
* Enqueue a message using a predefined Queue enum
*/
Mono<String> enqueue(Queue queue, Object message);
/**
* Enqueue a full RQ message to a specific queue
*/
Mono<String> enqueueJob(String queueName, RqMessage message);
/**
* Get the current size of a queue
*/
Mono<Integer> getQueueSize(String queueName);
}Benefits:
- Abstraction over queue implementation
- Easy to mock for testing
- Can be swapped with other implementations (Kafka, RabbitMQ)
Location: com.comet.opik.infrastructure.queues.Queue
public enum Queue {
OPTIMIZER_CLOUD("opik:optimizer-cloud", "opik_backend.rq_worker.process_optimizer_job");
private final String queueName;
private final String functionName; // Python function to call
}Usage:
queueProducer.enqueue(Queue.OPTIMIZER_CLOUD, myData);Benefits:
- Type-safe queue references
- Compile-time validation
- IDE autocomplete
- Queue name and function name coupled
Location: com.comet.opik.infrastructure.queues.RqMessage
public record RqMessage(
String id, // UUID
String func, // Python function name
List<Object> args, // Positional arguments
Map<String, Object> kwargs, // Keyword arguments
String description, // Job description
JobStatus status, // Job status (enum)
String origin, // Origin queue
Instant createdAt, // Creation timestamp
Instant enqueuedAt // Enqueued timestamp
) {
public static Builder builder() { ... }
}Benefits:
- Immutable by design
- Thread-safe
- Clear time semantics with
Instant - Type-safe status with enum
Location: com.comet.opik.infrastructure.queues.JobStatus
public enum JobStatus {
QUEUED, // Job has been queued but not started
STARTED, // Job is currently being executed
FINISHED, // Job finished successfully
FAILED; // Job failed during execution
}Location: com.comet.opik.infrastructure.redis.RqPublisher
Key methods:
class RqPublisher implements QueueProducer {
// Enqueue with type-safe Queue enum
public Mono<String> enqueue(Queue queue, Object message) {
RqMessage rqMessage = RqMessage.builder()
.func(queue.getFunctionName())
.args(List.of(message))
.origin(queue.toString())
.status(JobStatus.QUEUED)
.build();
return enqueueJob(queue.toString(), rqMessage);
}
// Low-level enqueue with full message control
public Mono<String> enqueueJob(String queueName, RqMessage message) {
String jobId = message.id();
String jobKey = "rq:job:" + jobId;
// Get TTL from configuration
Duration ttl = config.getQueues().getQueue(queueName)
.map(RqQueueConfig::getJobTTl)
.orElse(config.getQueues().getDefaultJobTtl());
// Store job data with TTL
return redisClient.getBucket(jobKey)
.set(message, ttl.toJavaDuration())
.then(redisClient.getQueue(queueName).offer(jobId));
}
}Location: apps/opik-python-backend/src/opik_backend/rq_worker.py
def process_optimizer_job(message: str):
"""Process an optimizer job from Java."""
logger.info(f"Processing optimizer job: {message}")
# Your processing logic here
result = {
"status": "success",
"message": f"Optimizer job processed: {message}",
"processed_by": "Python RQ Worker - Optimizer"
}
return result
def start_worker():
"""Start RQ worker listening on multiple queues."""
redis_conn = get_redis_connection()
queues = [
Queue("opik:hello_world_queue", connection=redis_conn),
Queue("opik:optimizer-cloud", connection=redis_conn),
]
worker = Worker(queues, connection=redis_conn)
worker.work()This section previously documented a zlib-based custom serializer and job class. The production path now uses RQ's native JSONSerializer and the default Job with plain JSON data. All custom serializer/job code has been removed.
The RQ worker includes comprehensive OpenTelemetry metrics for monitoring and observability. All metrics are automatically collected by the MetricsWorker class.
| Metric Name | Type | Description | Dimensions |
|---|---|---|---|
rq_worker.jobs.processed |
Counter | Total number of jobs processed (success + failure) | queue, function |
rq_worker.jobs.succeeded |
Counter | Number of successfully completed jobs | queue, function |
rq_worker.jobs.failed |
Counter | Number of failed jobs | queue, function, error_type |
| Metric Name | Type | Description | Unit | Dimensions |
|---|---|---|---|---|
rq_worker.job.processing_time |
Histogram | Time spent executing the job | milliseconds | queue, function |
rq_worker.job.queue_wait_time |
Histogram | Time job spent waiting in queue | milliseconds | queue, function |
rq_worker.job.total_time |
Histogram | Total time from creation to completion | milliseconds | queue, function |
All metrics include contextual dimensions for filtering and aggregation:
- queue: Queue name (e.g.,
opik:hello_world_queue,opik:optimizer-cloud) - function: Python function name (e.g.,
opik_backend.rq_worker.process_hello_world) - error_type: Exception class name (only for failed jobs, e.g.,
ValueError,ConnectionError)
The MetricsWorker extends RQ's standard Worker class and overrides perform_job() to collect metrics:
class MetricsWorker(Worker):
"""Custom RQ Worker that emits OpenTelemetry metrics."""
def perform_job(self, job, queue):
# Calculate queue wait time
if job.created_at and job.started_at:
queue_wait_ms = (job.started_at - job.created_at).total_seconds() * 1000
queue_wait_time_histogram.record(queue_wait_ms, {"queue": queue.name, "function": func_name})
# Execute job and measure processing time
result = super().perform_job(job, queue)
processing_time_ms = (time.time() - job_start_time) * 1000
# Record success metrics
jobs_processed_counter.add(1, {"queue": queue.name, "function": func_name})
jobs_succeeded_counter.add(1, {"queue": queue.name, "function": func_name})
processing_time_histogram.record(processing_time_ms, {"queue": queue.name, "function": func_name})Successful Job Processing:
rq_worker.jobs.processed{queue="opik:hello_world_queue", function="process_hello_world"} = 10
rq_worker.jobs.succeeded{queue="opik:hello_world_queue", function="process_hello_world"} = 10
rq_worker.job.processing_time{queue="opik:hello_world_queue", function="process_hello_world"} = [100ms, 102ms, 98ms, ...]
rq_worker.job.queue_wait_time{queue="opik:hello_world_queue", function="process_hello_world"} = [5ms, 3ms, 7ms, ...]
Failed Job Processing:
rq_worker.jobs.processed{queue="opik:optimizer-cloud", function="process_optimizer_job"} = 5
rq_worker.jobs.failed{queue="opik:optimizer-cloud", function="process_optimizer_job", error_type="ValueError"} = 1
from opentelemetry import metrics
from opentelemetry.sdk.metrics import MeterProvider
from opentelemetry.sdk.metrics.export import ConsoleMetricExporter, PeriodicExportingMetricReader
# Setup metric export
reader = PeriodicExportingMetricReader(ConsoleMetricExporter())
provider = MeterProvider(metric_readers=[reader])
metrics.set_meter_provider(provider)
# Metrics will be exported to console every 10 secondsThe metrics can be exported to various backends:
- Prometheus: Using
opentelemetry-exporter-prometheus - Jaeger: For distributed tracing
- Grafana: For visualization dashboards
- Cloud Providers: AWS CloudWatch, GCP Cloud Monitoring, Azure Monitor
-
Set up alerts for:
- High failure rate:
rq_worker.jobs.failed / rq_worker.jobs.processed > 0.05 - Long queue wait times:
rq_worker.job.queue_wait_time > 5000ms - Slow processing:
rq_worker.job.processing_time > 10000ms
- High failure rate:
-
Create dashboards showing:
- Jobs processed over time (throughput)
- Success vs failure rates
- Processing time percentiles (p50, p95, p99)
- Queue wait time trends
-
Track SLOs based on:
- 99.9% of jobs complete successfully
- 95% of jobs process within 1 second
- Queue wait time < 500ms for 99% of jobs
✅ Fully Implemented and Tested
- All 6 metrics defined and collecting data
- Dimensional data properly attached
- Integrated with RQ's job lifecycle
- No performance impact on job processing
- Ready for production observability platforms
queues:
# Enable/disable queue functionality
enabled: ${QUEUES_ENABLED:-true}
# Default TTL for all jobs (if not specified per-queue)
defaultJobTtl: ${QUEUES_DEFAULT_JOB_TTL:-1 day}
# Per-queue specific configurations
queues:
# Optimizer cloud queue
opik:optimizer-cloud:
jobTTl: ${OPTIMIZER_QUEUE_JOB_TTL:-1 day}
# Add more queue configs here
# opik:another-queue:
# jobTTl: 2 hours# Queue Configuration
QUEUES_ENABLED=true # Enable queue functionality
QUEUES_DEFAULT_JOB_TTL="1 day" # Default job TTL
OPTIMIZER_QUEUE_JOB_TTL="1 day" # Optimizer queue TTL
# Redis Connection
REDIS_URL="redis://:opik@localhost:6379/0"
REDIS_HOST=localhost
REDIS_PORT=6379
REDIS_DB=0
REDIS_PASSWORD=opik- Queue-specific TTL: Defined in
config.ymlunderqueues.queues.<queue-name>.jobTTl - Default TTL: Defined in
config.ymlunderqueues.defaultJobTtl - Fallback: If neither is set, uses 1 day
Duration ttl = config.getQueues()
.getQueue(queueName) // 1. Try queue-specific
.map(RqQueueConfig::getJobTTl)
.orElse(config.getQueues() // 2. Fallback to default
.getDefaultJobTtl());@Inject
private QueueProducer queueProducer;
public void sendOptimizationJob(String data) {
queueProducer.enqueue(Queue.OPTIMIZER_CLOUD, data)
.subscribe(
jobId -> log.info("Job enqueued: {}", jobId),
error -> log.error("Failed to enqueue", error)
);
}@Inject
private QueueProducer queueProducer;
public void sendCustomJob() {
RqMessage message = RqMessage.builder()
.func("opik_backend.rq_worker.process_custom_job")
.args(List.of("arg1", "arg2"))
.kwargs(Map.of("key1", "value1", "key2", "value2"))
.description("Custom job description")
.status(JobStatus.QUEUED)
.build();
queueProducer.enqueueJob("opik:custom-queue", message)
.subscribe(
jobId -> log.info("Custom job enqueued: {}", jobId),
error -> log.error("Failed to enqueue custom job", error)
);
}public Mono<Integer> getQueueDepth(Queue queue) {
return queueProducer.getQueueSize(queue.toString())
.doOnSuccess(size -> log.info("Queue {} size: {}", queue, size));
}public Mono<ProcessingResult> processWithQueue(String data) {
return validateData(data)
.flatMap(validated -> queueProducer.enqueue(Queue.OPTIMIZER_CLOUD, validated))
.flatMap(jobId -> waitForJobCompletion(jobId))
.map(result -> new ProcessingResult(result));
}File: com.comet.opik.infrastructure.queues.Queue
public enum Queue {
OPTIMIZER_CLOUD("opik:optimizer-cloud", "opik_backend.rq_worker.process_optimizer_job"),
// Add your new queue
MY_NEW_QUEUE("opik:my-new-queue", "opik_backend.rq_worker.process_my_new_job"),
;
}File: apps/opik-python-backend/src/opik_backend/rq_worker.py
def process_my_new_job(data: dict):
"""
Process my new job type.
Args:
data: The job data to process
Returns:
dict: Processing result
"""
logger.info(f"Processing my new job: {data}")
# Your processing logic
result = {
"status": "success",
"data": data,
"processed_at": datetime.now().isoformat()
}
logger.info("Job processed successfully")
return resultFile: apps/opik-python-backend/src/opik_backend/rq_worker.py
def start_worker():
redis_conn = get_redis_connection()
queues = [
Queue("opik:hello_world_queue", connection=redis_conn),
Queue("opik:optimizer-cloud", connection=redis_conn),
Queue("opik:my-new-queue", connection=redis_conn), # Add here
]
worker = Worker(queues, connection=redis_conn)
worker.work()File: apps/opik-backend/config.yml
queues:
queues:
opik:my-new-queue:
jobTTl: 2 hours # Custom TTL for this queue// In your service or resource
queueProducer.enqueue(Queue.MY_NEW_QUEUE, myData)
.subscribe(jobId -> log.info("Job enqueued: {}", jobId));# Clear Redis
redis-cli -a opik FLUSHDB
# Send test message
curl -X POST "http://localhost:8080/v1/internal/hello-world?message=test"
# Wait 2-3 seconds, then check status
redis-cli -a opik HGET rq:job:<job-id> status
# Expected: "finished"Test Results (2025-10-15):
✅ 10/10 messages sent (HTTP 200)
✅ 10/10 jobs finished successfully
✅ 0 failed jobs
✅ 100% success rate
Processing time: ~6 seconds for 10 jobs
Average: ~600ms per job (includes 500ms simulated processing)
Test Command:
# Clear and send 10 messages
redis-cli -a opik FLUSHDB
for i in {1..10}; do
curl -s -X POST "http://localhost:8080/v1/internal/hello-world?message=Test_${i}"
done
# Wait and check results
sleep 6
redis-cli -a opik KEYS 'rq:job:*' | wc -lVerified Features:
- ✅ Java creates RQ-compatible Redis HASH structures
- ✅ Plain JSON
data(UTF-8) with[func, null, args, kwargs] - ✅ RQ-native Redis keys (
rq:job:<id>,rq:queue:<queue>) - ✅
Job.fetch()and worker processing succeed with JSONSerializer - ✅ OpenTelemetry metrics infrastructure ready
@ExtendWith(MockitoExtension.class)
class MyServiceTest {
@Mock
private QueueProducer queueProducer;
@InjectMocks
private MyService myService;
@Test
void shouldEnqueueJobSuccessfully() {
// Given
String expectedJobId = "test-job-123";
when(queueProducer.enqueue(any(Queue.class), any()))
.thenReturn(Mono.just(expectedJobId));
// When
String result = myService.processData("test-data").block();
// Then
assertThat(result).isEqualTo(expectedJobId);
verify(queueProducer).enqueue(Queue.OPTIMIZER_CLOUD, "test-data");
}
}@Test
void shouldEnqueueAndProcessJob() throws InterruptedException {
// Given
String testMessage = "Integration test message";
// When - Enqueue job
String jobId = queueProducer.enqueue(Queue.OPTIMIZER_CLOUD, testMessage)
.block();
// Then - Verify job was enqueued
assertThat(jobId).isNotNull();
// Wait for Python worker to process (in real test, use polling or callbacks)
Thread.sleep(2000);
// Verify job was processed (check Redis or application state)
Integer queueSize = queueProducer.getQueueSize(Queue.OPTIMIZER_CLOUD.toString())
.block();
assertThat(queueSize).isZero();
}# Check job data (hash fields)
redis-cli HGETALL "rq:job:<job-id>"
# Check queue contents (RQ list)
redis-cli LRANGE "rq:queue:opik:optimizer-cloud" 0 -1
# Check queue length
redis-cli LLEN "rq:queue:opik:optimizer-cloud"
# Monitor Redis commands
redis-cli MONITORNone at the moment.
Symptom:
'utf-8' codec can't decode byte 0x9c in position 1: invalid start byte
Root cause:
datawas zlib-compressed; RQ restores jobs byHGETALLand attempts UTF-8 decoding of hash values before serializer runs.- The zlib header (
0x78 0x9c) triggered decode errors in that pre-serializer path.
Solution implemented:
- Switched
datato plain JSON (UTF-8) array:[func, null, args, kwargs]. - Use RQ's
JSONSerializerand defaultJobeverywhere (removed custom serializer/job). - Standardized Redis keys to RQ-native:
rq:job:<id>andrq:queue:<queue>. - Ensure a non-null
descriptionis written (prevents RQ logging issues).
Result:
- RQ worker processes Java-created jobs end-to-end reliably. Contract validated by tests and manual runs.
Symptoms: Jobs enqueued but never processed by Python worker
Checks:
# 1. Verify Python worker is running
ps aux | grep rq_worker
# 2. Check Redis queue
redis-cli -a opik LRANGE "opik:optimizer-cloud" 0 -1
# 3. Check job data exists
redis-cli -a opik KEYS "rq:job:*"
# 4. Check Python worker logs
tail -f /tmp/gunicorn.logSolutions:
- Ensure Python worker is started (via Gunicorn)
- Verify queue names match between Java and Python
- Check function names are correct
- Verify Redis connection in Python worker
Error: 'utf-8' codec can't decode byte 0x9c
Check:
# Verify job structure
redis-cli -a opik HGETALL "rq:job:<job-id>"
# Check if data field is binary
redis-cli -a opik HGET "rq:job:<job-id>" data | xxd | headSolution: This is the known limitation. See Current Limitations for potential workarounds.
Error: AttributeError: module 'opik_backend.rq_worker' has no attribute 'process_xxx'
Solution:
- Ensure function name in
Queueenum matches Python function name exactly - Check function is defined in
rq_worker.py - Verify Python module path is correct
Symptoms: Jobs disappear from Redis before being processed
Solution:
# Increase TTL in config.yml
queues:
defaultJobTtl: 7 days # Increase default
queues:
opik:my-queue:
jobTTl: 2 days # Or per-queueError: redis.exceptions.ConnectionError: Error connecting to Redis
Checks:
# Test Redis connectivity
redis-cli -h localhost -p 6379 PING
# Check Redis is running
docker ps | grep redis
# Test from Python
python -c "import redis; r = redis.Redis(); print(r.ping())"Solutions:
- Verify Redis is running
- Check
REDIS_HOSTandREDIS_PORTenvironment variables - Verify firewall rules allow Redis connection
- Check Redis authentication if configured
Error: TypeError: Object of type X is not JSON serializable
Solution:
- Ensure message data is JSON-serializable
- Convert complex objects to dictionaries
- Use strings, numbers, lists, and dictionaries only
// Bad - custom objects not serializable
MyCustomObject obj = new MyCustomObject();
queueProducer.enqueue(Queue.OPTIMIZER_CLOUD, obj); // ❌ Fails
// Good - use JSON-friendly types
Map<String, Object> data = Map.of(
"field1", obj.getField1(),
"field2", obj.getField2()
);
queueProducer.enqueue(Queue.OPTIMIZER_CLOUD, data); // ✅ WorksJava (config.yml):
logging:
loggers:
com.comet.opik.infrastructure.redis: DEBUG
com.comet.opik.infrastructure.queues: DEBUGPython:
logging.basicConfig(level=logging.DEBUG)redis-cli MONITOR | grep "opik:"# Get all job IDs
redis-cli KEYS "rq:job:*"
# Check specific job (hash)
redis-cli HGETALL "rq:job:<job-id>"
# Check queue (RQ list)
redis-cli LRANGE "rq:queue:opik:optimizer-cloud" 0 -1Decision: Use Java records instead of Lombok @Data classes
Reasons:
- Immutability: Records are immutable by default - thread-safe
- Less Boilerplate: No need for equals/hashCode/toString
- Modern Java: Idiomatic Java 16+ feature
- Clear Intent: Records signal immutable data carriers
Decision: Use java.time.Instant instead of Long (epoch millis/seconds)
Reasons:
- Type Safety: Strong typing prevents mixing seconds/millis
- Rich API: Built-in time manipulation methods
- ISO 8601: Standard serialization format
- Timezone Awareness: Better handling of time zones
- Clarity: Clear semantics - no guessing units
Decision: Use JobStatus enum instead of String
Reasons:
- Type Safety: Compile-time validation
- IDE Support: Autocomplete prevents typos
- Exhaustiveness: Switch statements warn if cases missing
- Documentation: Self-documenting valid states
Decision: Configure TTL at queue level, not per message
Reasons:
- Consistency: All jobs in a queue behave the same
- Separation of Concerns: Infrastructure config vs. message data
- Easier Management: Configure once per queue
- Flexibility: Different queues can have different policies
Decision: Create QueueProducer interface instead of using RqPublisher directly
Reasons:
- Dependency Inversion: Depend on abstraction, not implementation
- Testability: Easy to mock for unit tests
- Flexibility: Can swap implementations (Kafka, RabbitMQ)
- SOLID Principles: Interface Segregation Principle
Decision: Store full job data in bucket, only job ID in queue
Reasons:
- RQ Protocol: Required by Python RQ for job lifecycle management
- Separation: Queue for ordering, bucket for storage
- Efficiency: Only job IDs in queue (smaller memory footprint)
- Flexibility: Job data can be updated without touching queue
Original Structure:
infrastructure/rq/
├── RqPublisher.java (concrete class)
├── RqMessage.java (Lombok @Data)
├── RqQueueConfig.java (with factory methods)
└── JobStatus.java (not enum)
Issues:
- Tight coupling to concrete class
- Hardcoded TTL values
- String-based status (error-prone)
- Long timestamps (unit confusion)
- Complex factory methods
Changes:
- ✅ Converted
RqMessagefrom Lombok to record - ✅ Changed timestamps from
LongtoInstant - ✅ Created
JobStatusenum - ✅ Removed TTL from message, moved to queue config
Benefits:
- Immutability and thread safety
- Clear time semantics
- Type-safe status handling
- Consistent TTL per queue
Changes:
- ✅ Created
QueueProducerinterface - ✅ Created
Queueenum for type-safe queue definitions - ✅ Moved classes to proper packages (
queues/andredis/) - ✅ Added
QueuesConfigfor configuration - ✅ Integrated with Dropwizard config system
Benefits:
- Interface segregation
- Better package structure
- Configuration-driven design
- Easier to add new queues
└── infrastructure/
├── queues/ # Abstractions
│ ├── QueueProducer.java # Interface
│ ├── Queue.java # Enum
│ ├── RqMessage.java # Record
│ ├── RqQueueConfig.java # Config
│ └── JobStatus.java # Enum
├── redis/ # Implementation
│ └── RqPublisher.java # Concrete class
└── QueuesConfig.java # Configuration
-
SOLID Principles:
- Single Responsibility: Each class has one job
- Open/Closed: Open for extension (add queues), closed for modification
- Liskov Substitution:
RqPublishercan be substituted with anyQueueProducer - Interface Segregation: Small, focused
QueueProducerinterface - Dependency Inversion: Depend on
QueueProducer, notRqPublisher
-
DRY (Don't Repeat Yourself):
- Queue names and functions in one place (
Queueenum) - TTL logic centralized in configuration
- Queue names and functions in one place (
-
KISS (Keep It Simple):
- Simple interface with clear methods
- Minimal configuration required
- Sensible defaults
-
Immutability:
- Records are immutable
- Enums are constants
- Thread-safe by design
# Queue operations
RPUSH opik:optimizer-cloud <job-id> # Add job to queue
LPOP opik:optimizer-cloud # Remove job from queue
LLEN opik:optimizer-cloud # Get queue length
LRANGE opik:optimizer-cloud 0 -1 # View all jobs
# Job data operations
SET rq:job:<job-id> <json-data> # Store job data
GET rq:job:<job-id> # Get job data
DEL rq:job:<job-id> # Delete job data
TTL rq:job:<job-id> # Check TTL
# Monitoring
KEYS rq:job:* # List all jobs
KEYS opik:* # List all queues
MONITOR # Watch all commands# Start worker
python src/opik_backend/rq_worker.py
# Start with custom Redis
REDIS_HOST=custom-host REDIS_PORT=6380 python src/opik_backend/rq_worker.py
# View job status (using RQ CLI)
rq info --url redis://localhost:6379
# Empty queue
rq empty opik:optimizer-cloud --url redis://localhost:6379| Variable | Default | Description |
|---|---|---|
QUEUES_ENABLED |
true |
Enable queue functionality |
QUEUES_DEFAULT_JOB_TTL |
1 day |
Default job TTL |
OPTIMIZER_QUEUE_JOB_TTL |
1 day |
Optimizer queue job TTL |
REDIS_HOST |
localhost |
Redis host |
REDIS_PORT |
6379 |
Redis port |
REDIS_DB |
0 |
Redis database number |
REDIS_PASSWORD |
opik |
Redis password |
REDIS_URL |
redis://:opik@localhost:6379/0 |
Full Redis connection string |
For issues or questions:
- Check the Troubleshooting section
- Review logs in Java backend and Python worker
- Verify Redis connectivity and data
- Consult the Design Decisions for architecture rationale
Last Updated: 2025-10-15
Version: 2.0 (Post-Refactoring)