Guidelines for integrating new data series (FRED, CoinGecko) into background workers. Use this when expanding the dashboard's data coverage with new economic or crypto indicators.
This skill provides step-by-step guidance for adding new data indicators to the Cycle Navigator's background worker system. It ensures that new data integrations respect rate limits, caching logic, database constraints, and Celery task patterns.
Before adding a new data source, check your current usage:
FRED Rate Limits:
backend/tasks/fred_tasks.py to estimate cumulative daily callsCoinGecko Rate Limits:
backend/tasks/crypto_tasks.py for existing tasksNew data must follow this two-tier pattern:
External API (FRED/CoinGecko)
โ
PostgreSQL (source of truth)
โ
Redis Cache (24-hour TTL)
โ
API Response to Frontend
Add your task to the appropriate file:
backend/tasks/fred_tasks.pybackend/tasks/crypto_tasks.pyExample FRED Task:
from celery import shared_task
from backend.services.macro import fetch_and_store_fred_series
from backend.utils import get_redis_client
import logging
logger = logging.getLogger(__name__)
@shared_task(bind=True, max_retries=3)
def fetch_unemployment_rate(self):
"""Fetch monthly unemployment rate from FRED."""
series_id = "UNRATE"
redis_client = get_redis_client()
lock_key = f"fred:fetch:{series_id}"
try:
# Acquire distributed lock to prevent concurrent API calls
lock_acquired = redis_client.set(
lock_key,
"locked",
nx=True,
ex=300 # 5-minute lock timeout
)
if not lock_acquired:
logger.info(f"Lock held for {series_id}, skipping this run")
return
# Fetch from FRED API and store in PostgreSQL
data = fetch_and_store_fred_series(series_id)
# Update Redis cache (24-hour TTL)
cache_key = f"macro:fred:{series_id}"
redis_client.setex(cache_key, 86400, json.dumps(data))
logger.info(f"Successfully updated {series_id}")
except Exception as exc:
# Exponential backoff: 2s, 4s, 8s, 16s, 32s
retry_delay = 2 ** self.request.retries
self.retry(exc=exc, countdown=retry_delay)
finally:
# Always release the lock
redis_client.delete(lock_key)
Example Crypto Task:
@shared_task(bind=True, max_retries=3)
def fetch_bitcoin_market_cap(self):
"""Fetch Bitcoin market cap from CoinGecko."""
crypto_id = "bitcoin"
redis_client = get_redis_client()
lock_key = f"crypto:fetch:{crypto_id}"
try:
lock_acquired = redis_client.set(lock_key, "locked", nx=True, ex=60)
if not lock_acquired:
logger.info(f"Lock held for {crypto_id}, skipping")
return
# Fetch from CoinGecko
data = fetch_and_store_crypto_series(crypto_id)
# Cache with 24-hour TTL
cache_key = f"crypto:coingecko:{crypto_id}"
redis_client.setex(cache_key, 86400, json.dumps(data))
logger.info(f"Updated {crypto_id}")
except Exception as exc:
retry_delay = 2 ** self.request.retries
self.retry(exc=exc, countdown=retry_delay)
finally:
redis_client.delete(lock_key)
Edit backend/celery_app.py to schedule the new task:
from celery.schedules import crontab
# In app.conf.beat_schedule dictionary:
beat_schedule = {
# ... existing tasks ...
'fetch-unemployment-rate': {
'task': 'backend.tasks.fred_tasks.fetch_unemployment_rate',
'schedule': crontab(hour=0, minute=0), # Daily at midnight UTC
},
'fetch-bitcoin-market-cap': {
'task': 'backend.tasks.crypto_tasks.fetch_bitcoin_market_cap',
'schedule': crontab(minute='*/30'), # Every 30 minutes
},
}
If adding a new metric requires a new PostgreSQL table or column:
backend/models.py for existing schemaspython scripts/run_timescale_migrations.pyAdd a function to backend/services/macro.py or backend/services/crypto.py:
def fetch_and_store_fred_series(series_id: str) -> dict:
"""
Fetch from FRED API and store in PostgreSQL.
Returns: dict with latest data point
"""
# Call FRED API
fred = fredapi.FRED(api_key=settings.FRED_API_KEY)
data = fred.get(series_id)
# Store in PostgreSQL
# (Implementation depends on your ORM)
store_to_db(series_id, data)
return {
"series_id": series_id,
"latest_value": data.iloc[-1],
"timestamp": data.index[-1]
}
Add a route to backend/routers/macro.py or backend/routers/crypto.py:
@router.get("/unemployment-rate")
async def get_unemployment_rate():
"""Fetch cached unemployment rate data."""
redis_client = get_redis_client()
cached = redis_client.get("macro:fred:UNRATE")
if cached:
return json.loads(cached)
# Fallback to database if cache miss
return db_query_latest("UNRATE")
# Run a test fetch locally
podman-compose exec backend python -c "
from backend.tasks.fred_tasks import fetch_unemployment_rate
result = fetch_unemployment_rate()
print(result)
"
# Check data was written to database
podman-compose exec db psql -U postgres -d cycle_navigator -c "
SELECT * FROM macro_indicators WHERE series_id = 'UNRATE' LIMIT 5;
"
# Check cache entry
podman-compose exec redis redis-cli GET "macro:fred:UNRATE"
# Watch Celery worker logs
podman-compose logs -f celery-worker | grep fetch_unemployment_rate
You want to track monthly employment level (PAYEMS):
backend/tasks/fred_tasks.pybackend/celery_app.py: daily at 01:00 UTC (after FRED updates)podman-compose up -dYou want to track Ethereum daily volume:
backend/tasks/crypto_tasks.py/crypto/ethereum-volume endpointredis-cli MONITORDatabase Migrations:
Rate Limit Overages:
Redis Lock Failures:
redis:fetch:{series_id} naming convention consistentlyConcurrent Updates:
finally block to prevent stuck locksbackend/celery_app.py beat schedulepodman-composepsqlredis-cliCause: Task not registered in Celery app
Solution: Verify task is in backend/celery_app.py beat schedule and imported correctly
podman-compose exec backend celery -A backend.celery_app inspect active
Cause: Total daily API calls exceed service limits Solution: Reduce task frequency or remove lower-priority data
# Calculate current usage
grep "fetch_" backend/celery_app.py | wc -l
Cause: Lock held for too long or cache key mismatch Solution: Check lock TTL and cache key naming
redis-cli KEYS "fred:fetch:*" # Find stuck locks
redis-cli TTL "fred:fetch:UNRATE" # Check lock timeout
Cause: Duplicate data insertion or missing foreign keys
Solution: Review schema and check for UNIQUE or FOREIGN KEY constraints
podman-compose exec db psql -U postgres -d cycle_navigator -c "
SELECT constraint_name FROM information_schema.table_constraints
WHERE table_name = 'macro_indicators';
"
Cause: Task failing repeatedly and retrying with long delays Solution: Check API credentials, network connectivity, rate limits
# Review error logs
podman-compose logs celery-worker | grep "Retry"