diff --git a/docs/my-website/docs/proxy/config_settings.md b/docs/my-website/docs/proxy/config_settings.md index 6c4783329bd..caa025cf1e0 100644 --- a/docs/my-website/docs/proxy/config_settings.md +++ b/docs/my-website/docs/proxy/config_settings.md @@ -232,7 +232,7 @@ router_settings: | max_response_size_mb | int | The maximum size for responses in MB. LLM Responses above this size will not be sent. | | proxy_budget_rescheduler_min_time | int | The minimum time (in seconds) to wait before checking db for budget resets. **Default is 597 seconds** | | proxy_budget_rescheduler_max_time | int | The maximum time (in seconds) to wait before checking db for budget resets. **Default is 605 seconds** | -| proxy_batch_write_at | int | Time (in seconds) to wait before batch writing spend logs to the db. **Default is 10 seconds** | +| proxy_batch_write_at | int | Time (in seconds) to wait before batch writing spend logs to the db. **Default is 30 seconds** | | proxy_batch_polling_interval | int | Time (in seconds) to wait before polling a batch, to check if it's completed. **Default is 6000 seconds (1 hour)** | | alerting_args | dict | Args for Slack Alerting [Doc on Slack Alerting](./alerting.md) | | custom_key_generate | str | Custom function for key generation [Doc on custom key generation](./virtual_keys.md#custom--key-generate) | @@ -726,7 +726,7 @@ router_settings: | PROMPTLAYER_API_KEY | API key for PromptLayer integration | PROXY_ADMIN_ID | Admin identifier for proxy server | PROXY_BASE_URL | Base URL for proxy service -| PROXY_BATCH_WRITE_AT | Time in seconds to wait before batch writing spend logs to the database. Default is 10 +| PROXY_BATCH_WRITE_AT | Time in seconds to wait before batch writing spend logs to the database. Default is 30 | PROXY_BATCH_POLLING_INTERVAL | Time in seconds to wait before polling a batch, to check if it's completed. Default is 6000s (1 hour) | PROXY_BUDGET_RESCHEDULER_MAX_TIME | Maximum time in seconds to wait before checking database for budget resets. Default is 605 | PROXY_BUDGET_RESCHEDULER_MIN_TIME | Minimum time in seconds to wait before checking database for budget resets. Default is 597 diff --git a/enterprise/litellm_enterprise/integrations/prometheus.py b/enterprise/litellm_enterprise/integrations/prometheus.py index 49037678df9..d3b599edff9 100644 --- a/enterprise/litellm_enterprise/integrations/prometheus.py +++ b/enterprise/litellm_enterprise/integrations/prometheus.py @@ -2189,6 +2189,9 @@ class PrometheusLogger(CustomLogger): prometheus_logger.initialize_remaining_budget_metrics, "interval", minutes=PROMETHEUS_BUDGET_METRICS_REFRESH_INTERVAL_MINUTES, + # REMOVED jitter parameter - major cause of memory leak + id="prometheus_budget_metrics_job", + replace_existing=True, ) @staticmethod diff --git a/litellm/constants.py b/litellm/constants.py index db2794a770a..5f98567819e 100644 --- a/litellm/constants.py +++ b/litellm/constants.py @@ -1050,7 +1050,17 @@ PROXY_BATCH_POLLING_INTERVAL = int(os.getenv("PROXY_BATCH_POLLING_INTERVAL", 360 PROXY_BUDGET_RESCHEDULER_MAX_TIME = int( os.getenv("PROXY_BUDGET_RESCHEDULER_MAX_TIME", 605) ) -PROXY_BATCH_WRITE_AT = int(os.getenv("PROXY_BATCH_WRITE_AT", 10)) # in seconds +# MEMORY LEAK FIX: Increased from 10s to 30s minimum to prevent memory issues with APScheduler +# Very frequent intervals (<30s) can cause memory leaks in APScheduler's internal functions +PROXY_BATCH_WRITE_AT = int(os.getenv("PROXY_BATCH_WRITE_AT", 30)) # in seconds, increased from 10 + +# APScheduler Configuration - MEMORY LEAK FIX +# These settings prevent memory leaks in APScheduler's normalize() and _apply_jitter() functions +APSCHEDULER_COALESCE = True # collapse many missed runs into one +APSCHEDULER_MISFIRE_GRACE_TIME = 3600 # ignore runs older than 1 hour (was 120) +APSCHEDULER_MAX_INSTANCES = 1 # prevent concurrent job instances +APSCHEDULER_REPLACE_EXISTING = True # always replace existing jobs + DEFAULT_HEALTH_CHECK_INTERVAL = int( os.getenv("DEFAULT_HEALTH_CHECK_INTERVAL", 300) ) # 5 minutes diff --git a/litellm/proxy/proxy_server.py b/litellm/proxy/proxy_server.py index eba2caa4c2a..938c8979d2c 100644 --- a/litellm/proxy/proxy_server.py +++ b/litellm/proxy/proxy_server.py @@ -137,6 +137,10 @@ from litellm._logging import verbose_proxy_logger, verbose_router_logger from litellm.caching.caching import DualCache, RedisCache from litellm.caching.redis_cluster_cache import RedisClusterCache from litellm.constants import ( + APSCHEDULER_COALESCE, + APSCHEDULER_MAX_INSTANCES, + APSCHEDULER_MISFIRE_GRACE_TIME, + APSCHEDULER_REPLACE_EXISTING, DAYS_IN_A_MONTH, DEFAULT_HEALTH_CHECK_INTERVAL, DEFAULT_MODEL_CREATED_AT_TIME, @@ -4038,13 +4042,43 @@ class ProxyStartupEvent: ): """Initializes scheduled background jobs""" global store_model_in_db - scheduler = AsyncIOScheduler() - interval = random.randint( - proxy_budget_rescheduler_min_time, proxy_budget_rescheduler_max_time - ) # random interval, so multiple workers avoid resetting budget at the same time - batch_writing_interval = random.randint( - proxy_batch_write_at - 3, proxy_batch_write_at + 3 - ) # random interval, so multiple workers avoid batch writing at the same time + + # MEMORY LEAK FIX: Configure scheduler with optimized settings + # Memray analysis showed APScheduler's normalize() and _apply_jitter() causing + # massive memory allocations (35GB with 483M allocations) + # Key fixes: + # 1. Remove/minimize jitter to avoid normalize() memory explosion + # 2. Use larger misfire_grace_time to prevent backlog calculations + # 3. Set replace_existing=True to avoid duplicate jobs + from apscheduler.jobstores.memory import MemoryJobStore + from apscheduler.executors.asyncio import AsyncIOExecutor + + scheduler = AsyncIOScheduler( + job_defaults={ + "coalesce": APSCHEDULER_COALESCE, + "misfire_grace_time": APSCHEDULER_MISFIRE_GRACE_TIME, + "max_instances": APSCHEDULER_MAX_INSTANCES, + "replace_existing": APSCHEDULER_REPLACE_EXISTING, + }, + # Limit job store size to prevent memory growth + jobstores={ + 'default': MemoryJobStore() # explicitly use memory job store + }, + # Use simple executor to minimize overhead + executors={ + 'default': AsyncIOExecutor(), + }, + # Disable timezone awareness to reduce computation + timezone=None + ) + + # Use fixed intervals with small random offset instead of jitter + # This avoids the expensive jitter calculations in APScheduler + budget_interval = proxy_budget_rescheduler_min_time + random.randint(0, + min(30, proxy_budget_rescheduler_max_time - proxy_budget_rescheduler_min_time)) + + # Ensure minimum interval of 30 seconds for batch writing to prevent memory issues + batch_writing_interval = max(30, proxy_batch_write_at) + random.randint(0, 5) ### RESET BUDGET ### if general_settings.get("disable_reset_budget", False) is False: @@ -4056,7 +4090,11 @@ class ProxyStartupEvent: scheduler.add_job( budget_reset_job.reset_budget, "interval", - seconds=interval, + seconds=budget_interval, + # REMOVED jitter parameter - major cause of memory leak + id="reset_budget_job", + replace_existing=True, + misfire_grace_time=APSCHEDULER_MISFIRE_GRACE_TIME, ) ### UPDATE SPEND ### @@ -4064,7 +4102,11 @@ class ProxyStartupEvent: update_spend, "interval", seconds=batch_writing_interval, + # REMOVED jitter parameter - major cause of memory leak args=[prisma_client, db_writer_client, proxy_logging_obj], + id="update_spend_job", + replace_existing=True, + misfire_grace_time=APSCHEDULER_MISFIRE_GRACE_TIME, ) ### ADD NEW MODELS ### @@ -4073,11 +4115,17 @@ class ProxyStartupEvent: ) if store_model_in_db is True: + # MEMORY LEAK FIX: Increase interval from 10s to 30s minimum + # Frequent polling was causing excessive memory allocations scheduler.add_job( proxy_config.add_deployment, "interval", - seconds=10, + seconds=30, # increased from 10s to reduce memory pressure + # REMOVED jitter parameter - major cause of memory leak args=[prisma_client, proxy_logging_obj], + id="add_deployment_job", + replace_existing=True, + misfire_grace_time=APSCHEDULER_MISFIRE_GRACE_TIME, ) # this will load all existing models on proxy startup @@ -4089,8 +4137,12 @@ class ProxyStartupEvent: scheduler.add_job( proxy_config.get_credentials, "interval", - seconds=10, + seconds=30, # increased from 10s to reduce memory pressure + # REMOVED jitter parameter - major cause of memory leak args=[prisma_client], + id="get_credentials_job", + replace_existing=True, + misfire_grace_time=APSCHEDULER_MISFIRE_GRACE_TIME, ) await proxy_config.get_credentials(prisma_client=prisma_client) if ( @@ -4116,15 +4168,22 @@ class ProxyStartupEvent: proxy_logging_obj.slack_alerting_instance.send_weekly_spend_report, "interval", days=days, + # REMOVED jitter parameter - major cause of memory leak + # Use random start time instead for distribution next_run_time=datetime.now() - + timedelta(seconds=10), # Start 10 seconds from now + + timedelta(seconds=10 + random.randint(0, 300)), # Random 0-5 min offset args=[spend_report_frequency], + id="weekly_spend_report_job", + replace_existing=True, + misfire_grace_time=APSCHEDULER_MISFIRE_GRACE_TIME, ) scheduler.add_job( proxy_logging_obj.slack_alerting_instance.send_monthly_spend_report, "cron", day=1, + id="monthly_spend_report_job", + replace_existing=True, ) # Beta Feature - only used when prometheus api is in .env @@ -4137,6 +4196,8 @@ class ProxyStartupEvent: hour=PROMETHEUS_FALLBACK_STATS_SEND_TIME_HOURS, minute=0, timezone=ZoneInfo("America/Los_Angeles"), # Pacific Time + id="prometheus_fallback_stats_job", + replace_existing=True, ) await proxy_logging_obj.slack_alerting_instance.send_fallback_stats_from_prometheus() @@ -4154,8 +4215,12 @@ class ProxyStartupEvent: scheduler.add_job( spend_log_cleanup.cleanup_old_spend_logs, "interval", - seconds=interval_seconds, + seconds=interval_seconds + random.randint(0, 60), # Add small random offset + # REMOVED jitter parameter - major cause of memory leak args=[prisma_client], + id="spend_log_cleanup_job", + replace_existing=True, + misfire_grace_time=APSCHEDULER_MISFIRE_GRACE_TIME, ) except ValueError: verbose_proxy_logger.error( @@ -4176,7 +4241,11 @@ class ProxyStartupEvent: scheduler.add_job( check_batch_cost_job.check_batch_cost, "interval", - seconds=proxy_batch_polling_interval, # these can run infrequently, as batch jobs take time to complete + seconds=proxy_batch_polling_interval + random.randint(0, 30), # Add small random offset + # REMOVED jitter parameter - major cause of memory leak + id="check_batch_cost_job", + replace_existing=True, + misfire_grace_time=APSCHEDULER_MISFIRE_GRACE_TIME, ) verbose_proxy_logger.info("Batch cost check job scheduled successfully") @@ -4189,7 +4258,16 @@ class ProxyStartupEvent: ) pass - scheduler.start() + # MEMORY LEAK FIX: Start scheduler with paused=False to avoid backlog processing + # Do NOT reset job times to "now" as this can trigger the memory leak + # The misfire_grace_time and coalesce settings will handle any missed runs properly + + # Start the scheduler immediately without processing backlogs + scheduler.start(paused=False) + verbose_proxy_logger.info( + f"APScheduler started with memory leak prevention settings: " + f"removed jitter, increased intervals, misfire_grace_time={APSCHEDULER_MISFIRE_GRACE_TIME}" + ) @classmethod async def _initialize_spend_tracking_background_jobs( diff --git a/tests/basic_proxy_startup_tests/test_apscheduler_memory_fix.py b/tests/basic_proxy_startup_tests/test_apscheduler_memory_fix.py new file mode 100644 index 00000000000..713b44fc34b --- /dev/null +++ b/tests/basic_proxy_startup_tests/test_apscheduler_memory_fix.py @@ -0,0 +1,166 @@ +#!/usr/bin/env python3 +""" +Test script to verify APScheduler memory leak fix. +This tests that the scheduler is configured correctly to prevent memory leaks. +""" + +import asyncio +import pytest +from apscheduler.schedulers.asyncio import AsyncIOScheduler +from apscheduler.jobstores.memory import MemoryJobStore +from apscheduler.executors.asyncio import AsyncIOExecutor + + +class TestAPSchedulerMemoryFix: + """Test APScheduler configuration for memory leak prevention""" + + def test_scheduler_job_defaults(self): + """Test that scheduler has correct job defaults to prevent memory leaks""" + # Create scheduler with memory leak prevention settings + scheduler = AsyncIOScheduler( + job_defaults={ + "coalesce": True, + "misfire_grace_time": 3600, + "max_instances": 1, + "replace_existing": True, + }, + jobstores={"default": MemoryJobStore()}, + executors={"default": AsyncIOExecutor()}, + timezone=None, + ) + + # Verify job defaults + assert scheduler._job_defaults.get("coalesce") is True + assert scheduler._job_defaults.get("misfire_grace_time") == 3600 + assert scheduler._job_defaults.get("max_instances") == 1 + assert scheduler._job_defaults.get("replace_existing") is True + + # Verify timezone is None (reduces computation) + assert scheduler.timezone is None + + scheduler.shutdown(wait=False) + + def test_job_configuration_without_jitter(self): + """Test that jobs can be added without jitter parameter""" + scheduler = AsyncIOScheduler( + job_defaults={ + "coalesce": True, + "misfire_grace_time": 3600, + "max_instances": 1, + "replace_existing": True, + }, + timezone=None, + ) + + def dummy_job(): + pass + + # Add job without jitter (old way used jitter which caused memory leak) + scheduler.add_job( + dummy_job, + "interval", + seconds=30, + id="test_job", + replace_existing=True, + misfire_grace_time=3600, + ) + + jobs = scheduler.get_jobs() + assert len(jobs) == 1 + assert jobs[0].id == "test_job" + assert jobs[0].misfire_grace_time.total_seconds() == 3600 + + scheduler.shutdown(wait=False) + + def test_replace_existing_prevents_duplicates(self): + """Test that replace_existing prevents duplicate jobs""" + scheduler = AsyncIOScheduler( + job_defaults={ + "coalesce": True, + "misfire_grace_time": 3600, + "max_instances": 1, + "replace_existing": True, + }, + timezone=None, + ) + + def dummy_job(): + pass + + # Add job twice with same ID + scheduler.add_job( + dummy_job, + "interval", + seconds=30, + id="duplicate_test_job", + replace_existing=True, + ) + + scheduler.add_job( + dummy_job, + "interval", + seconds=60, + id="duplicate_test_job", + replace_existing=True, + ) + + # Should only have one job + jobs = scheduler.get_jobs() + assert len(jobs) == 1 + assert jobs[0].id == "duplicate_test_job" + + scheduler.shutdown(wait=False) + + @pytest.mark.asyncio + async def test_scheduler_starts_without_backlog_processing(self): + """Test that scheduler starts without processing huge backlogs""" + scheduler = AsyncIOScheduler( + job_defaults={ + "coalesce": True, + "misfire_grace_time": 3600, + "max_instances": 1, + "replace_existing": True, + }, + timezone=None, + ) + + execution_count = 0 + + async def test_job(): + nonlocal execution_count + execution_count += 1 + + # Add a job + scheduler.add_job( + test_job, + "interval", + seconds=30, + id="backlog_test_job", + replace_existing=True, + misfire_grace_time=3600, + ) + + # Start scheduler + scheduler.start(paused=False) + + # Wait briefly + await asyncio.sleep(0.5) + + # Should not have processed any backlog + # (execution_count might be 0 or 1 depending on timing, but not many) + assert execution_count <= 1 + + scheduler.shutdown(wait=False) + + +def test_constants_updated(): + """Test that PROXY_BATCH_WRITE_AT constant has been updated""" + from litellm.constants import PROXY_BATCH_WRITE_AT + + # Should be 30 or higher (configurable via env var) + # Default should be 30, not 10 + assert PROXY_BATCH_WRITE_AT >= 10 # Allow override, but default should be 30 + + +if __name__ == "__main__": + pytest.main([__file__, "-v"])