Rupesh Raut
04/20/2026, 7:47 AMMarvin
04/20/2026, 7:47 AMMarvin
04/20/2026, 7:47 AMRupesh Raut
04/20/2026, 7:48 AMMarvin
04/20/2026, 7:48 AMMarvin
04/20/2026, 7:48 AMprefect.yaml
I can:
- Research a Prefect topic against the source code
- Verify CLI commands and flags
- Search GitHub issues/discussions for similar reports
Tell me which path you want, and I’ll run the right tools.Rupesh Raut
04/20/2026, 7:48 AMMarvin
04/20/2026, 7:48 AMMarvin
04/20/2026, 7:49 AMRupesh Raut
04/20/2026, 7:50 AMMarvin
04/20/2026, 7:50 AMMarvin
04/20/2026, 7:54 AMexport PREFECT_LOG_LEVEL=DEBUG
# restart your workers and server services, then reproduce for ~10–20 minutes
- In server service logs (TaskRunRecorder), look for:
- “IntegrityError” or “Retrying task_run upsert…”
- “Dropping event … after … failed attempts”
- Long gaps between “flushed N task runs” messages
- In worker logs, look for:
- “Failed to emit heartbeat” or unusually long delays between task finish and “state submitted” messages
4) Increase TaskRunRecorder/EventPersister throughput on the server
- The defaults are conservative and can lag under concurrency. Scale batching on the server service(s):
# TaskRunRecorder (no legacy aliases)
export PREFECT_SERVER_SERVICES_TASK_RUN_RECORDER_BATCH_SIZE=50
export PREFECT_SERVER_SERVICES_TASK_RUN_RECORDER_READ_BATCH_SIZE=10
export PREFECT_SERVER_SERVICES_TASK_RUN_RECORDER_FLUSH_INTERVAL=1
# EventPersister (supports legacy aliases but prefer these)
export PREFECT_SERVER_SERVICES_EVENT_PERSISTER_BATCH_SIZE=100
export PREFECT_SERVER_SERVICES_EVENT_PERSISTER_READ_BATCH_SIZE=10
export PREFECT_SERVER_SERVICES_EVENT_PERSISTER_FLUSH_INTERVAL=2
- For very high concurrency (100+ concurrent tasks), consider:
export PREFECT_SERVER_SERVICES_TASK_RUN_RECORDER_BATCH_SIZE=200
export PREFECT_SERVER_SERVICES_TASK_RUN_RECORDER_READ_BATCH_SIZE=50
export PREFECT_SERVER_SERVICES_TASK_RUN_RECORDER_FLUSH_INTERVAL=0.5
5) If CPU-bound workloads, relax heartbeat frequency a bit
- This reduces scheduling pressure on the heartbeat/state-emission thread:
export PREFECT_FLOWS_HEARTBEAT_FREQUENCY=60
- Note: Prefect enforces a 30s minimum; 60–120s is a good target for busy nodes. There is no separate “task heartbeat” setting in 3.x.
6) Check Postgres for stuck or delayed state persistence
- Look for long-running RUNNING states and state timestamp anomalies:
```
-- Task runs stuck in RUNNING for >10 minutes
SELECT id, flow_run_id, task_key, dynamic_key, state_type, state_name, state_timestamp, updated
FROM task_run
WHERE state_type = 'RUNNING' AND updated < NOW() - INTERVAL '10 minutes'Marvin
04/20/2026, 7:54 AM- Check DB logs for deadlocks or unique constraint violations mentioning task_run/task_run_state.
7) Double-check version alignment and broker configuration
- Server services and workers should:
- Run the same Prefect version (ideally latest 3.6.x).
- Share the same Redis URL and credentials.
- Have adequate DB connection pools on the server:
export PREFECT_SERVER_DATABASE_CONNECTION_POOL_MIN_SIZE=5
export PREFECT_SERVER_DATABASE_CONNECTION_POOL_MAX_SIZE=20
```
Why this happens after your changes
- The upgrade introduced more robust but also more active event recording; combined with your shift to multiple workers + Redis, you likely increased event volume and concurrency. With default batching (batch_size=1, flush_interval=5s), the recorder can fall behind. Under concurrency, brief DB conflicts can occur; without retries and larger batches, terminal states may be delayed long enough to stall downstream tasks.
If you want, I can help triage live with a focused checklist. Helpful bits to share:
- Confirm self-hosted vs Cloud, and exact Prefect version of:
- workers
- server API
- server services image(s)
- Whether workers are CPU-bound during these runs
- Redis URL config locations for both workers and server services
- A short log excerpt around a single “stuck” task run from:
- The worker that ran it (DEBUG level)
- The server TaskRunRecorder (DEBUG level)
- Output of the SQL above for stuck RUNNING tasks
Optional mitigations while you tune/upgrade
- Temporarily lower concurrency or number of workers to reduce event pressure.
- Restart the TaskRunRecorder service if you see it stuck with a large queue (after bumping batch sizes).
- Upgrade to the latest 3.6.x first; if the problem persists, we can instrument the recorder and broker to pinpoint drops vs. delays.
I’m happy to tailor exact env var manifests (Kubernetes, Docker Compose, systemd) if you share how you’re running the server services and workers.Rupesh Raut
04/20/2026, 8:05 AMPREFECT_API_MAX_FLOW_RUN_GRAPH_ARTIFACTS='10000' (from defaults)
PREFECT_API_MAX_FLOW_RUN_GRAPH_NODES='10000' (from defaults)
PREFECT_API_REQUEST_TIMEOUT='60.0' (from defaults)
PREFECT_API_SERVICES_CANCELLATION_CLEANUP_ENABLED='True' (from defaults)
PREFECT_API_SERVICES_CANCELLATION_CLEANUP_LOOP_SECONDS='20.0' (from defaults)
PREFECT_API_SERVICES_EVENT_LOGGER_ENABLED='False' (from defaults)
PREFECT_API_SERVICES_EVENT_PERSISTER_BATCH_SIZE='20' (from defaults)
PREFECT_API_SERVICES_EVENT_PERSISTER_ENABLED='True' (from defaults)
PREFECT_API_SERVICES_EVENT_PERSISTER_FLUSH_INTERVAL='5.0' (from defaults)
PREFECT_API_SERVICES_EVENT_PERSISTER_READ_BATCH_SIZE='1' (from defaults)
PREFECT_API_SERVICES_FOREMAN_DEPLOYMENT_LAST_POLLED_TIMEOUT_SECONDS='60' (from defaults)
PREFECT_API_SERVICES_FOREMAN_ENABLED='True' (from defaults)
PREFECT_API_SERVICES_FOREMAN_FALLBACK_HEARTBEAT_INTERVAL_SECONDS='30' (from defaults)
PREFECT_API_SERVICES_FOREMAN_INACTIVITY_HEARTBEAT_MULTIPLE='3' (from defaults)
PREFECT_API_SERVICES_FOREMAN_LOOP_SECONDS='15.0' (from defaults)
PREFECT_API_SERVICES_FOREMAN_WORK_QUEUE_LAST_POLLED_TIMEOUT_SECONDS='60' (from defaults)
PREFECT_API_SERVICES_LATE_RUNS_AFTER_SECONDS='0:00:15' (from defaults)
PREFECT_API_SERVICES_LATE_RUNS_ENABLED='True' (from defaults)
PREFECT_API_SERVICES_LATE_RUNS_LOOP_SECONDS='5.0' (from defaults)
PREFECT_API_SERVICES_PAUSE_EXPIRATIONS_ENABLED='True' (from defaults)
PREFECT_API_SERVICES_PAUSE_EXPIRATIONS_LOOP_SECONDS='5.0' (from defaults)
PREFECT_API_SERVICES_SCHEDULER_DEPLOYMENT_BATCH_SIZE='100' (from defaults)
PREFECT_API_SERVICES_SCHEDULER_ENABLED='True' (from defaults)
PREFECT_API_SERVICES_SCHEDULER_INSERT_BATCH_SIZE='500' (from defaults)
PREFECT_API_SERVICES_SCHEDULER_LOOP_SECONDS='60.0' (from defaults)
PREFECT_API_SERVICES_SCHEDULER_MAX_RUNS='100' (from defaults)
PREFECT_API_SERVICES_SCHEDULER_MAX_SCHEDULED_TIME='100 days, 0:00:00' (from defaults)
PREFECT_API_SERVICES_SCHEDULER_MIN_RUNS='3' (from defaults)
PREFECT_API_SERVICES_SCHEDULER_MIN_SCHEDULED_TIME='1:00:00' (from defaults)
PREFECT_API_SERVICES_TASK_RUN_RECORDER_ENABLED='True' (from defaults)
PREFECT_API_SERVICES_TRIGGERS_ENABLED='True' (from defaults)
PREFECT_API_SSL_CERT_FILE='None' (from defaults)
PREFECT_API_TASK_CACHE_KEY_MAX_LENGTH='2000' (from defaults)
PREFECT_API_TLS_INSECURE_SKIP_VERIFY='False' (from defaults)
PREFECT_MESSAGING_BROKER='prefect_redis.messaging' (from profile)
PREFECT_MESSAGING_CACHE='prefect_redis.messaging' (from profile)
PREFECT_PROFILES_PATH='/home/infra/.prefect/profiles.toml' (from defaults)
PREFECT_RESULTS_DEFAULT_SERIALIZER='pickle' (from defaults)
PREFECT_RESULTS_PERSIST_BY_DEFAULT='true' (from profile)
PREFECT_RUNNER_CRASH_ON_CANCELLATION_FAILURE='False' (from defaults)
PREFECT_RUNNER_HEARTBEAT_FREQUENCY='180' (from defaults)
PREFECT_RUNNER_POLL_FREQUENCY='10' (from defaults)
PREFECT_RUNNER_PROCESS_LIMIT='5' (from defaults)
PREFECT_RUNNER_SERVER_ENABLE='False' (from defaults)
PREFECT_RUNNER_SERVER_HOST='localhost' (from defaults)
PREFECT_RUNNER_SERVER_LOG_LEVEL='ERROR' (from defaults)
PREFECT_RUNNER_SERVER_MISSED_POLLS_TOLERANCE='2' (from defaults)
PREFECT_RUNNER_SERVER_PORT='8080' (from defaults)
PREFECT_SERVER_ALLOW_EPHEMERAL_MODE='False' (from defaults)
PREFECT_SERVER_ANALYTICS_ENABLED='True' (from defaults)
PREFECT_SERVER_API_AUTH_STRING='********' (from defaults)
PREFECT_SERVER_API_BASE_PATH='None' (from defaults)
PREFECT_SERVER_API_HOST='0.0.0.0' (from profile)
PREFECT_SERVER_API_KEEPALIVE_TIMEOUT='5' (from defaults)
PREFECT_SERVER_API_PORT='4200' (from profile)
PREFECT_SERVER_CONCURRENCY_INITIAL_DEPLOYMENT_LEASE_DURATION='300.0' (from defaults)
PREFECT_SERVER_CONCURRENCY_LEASE_STORAGE='prefect_redis.lease_storage' (from profile)
PREFECT_SERVER_CONCURRENCY_MAXIMUM_CONCURRENCY_SLOT_WAIT_SECONDS='30.0' (from defaults)
PREFECT_SERVER_CORS_ALLOWED_HEADERS='*' (from defaults)
PREFECT_SERVER_CORS_ALLOWED_METHODS='*' (from defaults)
PREFECT_SERVER_CORS_ALLOWED_ORIGINS='*' (from defaults)
PREFECT_SERVER_CSRF_PROTECTION_ENABLED='False' (from defaults)
PREFECT_SERVER_CSRF_TOKEN_EXPIRATION='1:00:00' (from defaults)
PREFECT_SERVER_DATABASE_SQLALCHEMY_CONNECT_ARGS_APPLICATION_NAME='None' (from defaults)
PREFECT_SERVER_DATABASE_SQLALCHEMY_CONNECT_ARGS_PREPARED_STATEMENT_CACHE_SIZE='None' (from defaults)
PREFECT_SERVER_DATABASE_SQLALCHEMY_CONNECT_ARGS_SEARCH_PATH='None' (from defaults)
PREFECT_SERVER_DATABASE_SQLALCHEMY_CONNECT_ARGS_STATEMENT_CACHE_SIZE='None' (from defaults)
PREFECT_SERVER_DATABASE_SQLALCHEMY_CONNECT_ARGS_TLS_CA_FILE='None' (from defaults)
PREFECT_SERVER_DATABASE_SQLALCHEMY_CONNECT_ARGS_TLS_CERT_FILE='None' (from defaults)
PREFECT_SERVER_DATABASE_SQLALCHEMY_CONNECT_ARGS_TLS_CHECK_HOSTNAME='True' (from defaults)
PREFECT_SERVER_DATABASE_SQLALCHEMY_CONNECT_ARGS_TLS_ENABLED='False' (from defaults)
PREFECT_SERVER_DATABASE_SQLALCHEMY_CONNECT_ARGS_TLS_KEY_FILE='None' (from defaults)
PREFECT_SERVER_DATABASE_SQLALCHEMY_POOL_RECYCLE='3600' (from defaults)
PREFECT_SERVER_DATABASE_SQLALCHEMY_POOL_TIMEOUT='30.0' (from defaults)
PREFECT_SERVER_DOCKET_NAME='prefect-server' (from defaults)
PREFECT_SERVER_DOCKET_URL='memory://' (from defaults)
PREFECT_SERVER_EPHEMERAL_STARTUP_TIMEOUT_SECONDS='20' (from defaults)
PREFECT_SERVER_EVENTS_CAUSAL_ORDERING='prefect_redis.ordering' (from profile)
PREFECT_SERVER_EVENTS_MAXIMUM_EVENT_NAME_LENGTH='1024' (from defaults)
PREFECT_SERVER_LOGS_STREAM_OUT_ENABLED='False' (from defaults)
PREFECT_SERVER_LOGS_STREAM_PUBLISHING_ENABLED='False' (from defaults)
PREFECT_SERVER_SERVICES_DB_VACUUM_BATCH_SIZE='200' (from defaults)
PREFECT_SERVER_SERVICES_DB_VACUUM_ENABLED='{'events'}' (from defaults)
PREFECT_SERVER_SERVICES_DB_VACUUM_LOOP_SECONDS='3600.0' (from defaults)
PREFECT_SERVER_SERVICES_DB_VACUUM_RETENTION_PERIOD='90 days, 0:00:00' (from defaults)
PREFECT_SERVER_SERVICES_EVENT_PERSISTER_MAX_FLUSH_RETRIES='5' (from defaults)
PREFECT_SERVER_SERVICES_EVENT_PERSISTER_QUEUE_MAX_SIZE='50000' (from defaults)
PREFECT_SERVER_SERVICES_REPOSSESSOR_ENABLED='True' (from defaults)
PREFECT_SERVER_SERVICES_REPOSSESSOR_LOOP_SECONDS='15.0' (from defaults)
PREFECT_SERVER_SERVICES_SCHEDULER_RECENT_DEPLOYMENTS_LOOP_SECONDS='5.0' (from defaults)
PREFECT_SERVER_SERVICES_TASK_RUN_RECORDER_BATCH_SIZE='1' (from defaults)
PREFECT_SERVER_SERVICES_TASK_RUN_RECORDER_FLUSH_INTERVAL='5.0' (from defaults)
PREFECT_SERVER_SERVICES_TASK_RUN_RECORDER_READ_BATCH_SIZE='1' (from defaults)
PREFECT_SERVER_SERVICES_TRIGGERS_PG_NOTIFY_HEARTBEAT_INTERVAL_SECONDS='5' (from defaults)
PREFECT_SERVER_SERVICES_TRIGGERS_PG_NOTIFY_RECONNECT_INTERVAL_SECONDS='10' (from defaults)
PREFECT_SERVER_SERVICES_TRIGGERS_READ_BATCH_SIZE='1' (from defaults)
PREFECT_SERVER_UI_SHOW_PROMOTIONAL_CONTENT='True' (from defaults)
PREFECT_SERVER_UI_V2_ENABLED='False' (from defaults)
PREFECT_SILENCE_API_URL_MISCONFIGURATION='False' (from defaults)
PREFECT_SQLALCHEMY_MAX_OVERFLOW='10' (from defaults)
PREFECT_SQLALCHEMY_POOL_SIZE='5' (from defaults)
PREFECT_TASKS_DEFAULT_NO_CACHE='False' (from defaults)
PREFECT_TASKS_DEFAULT_PERSIST_RESULT='None' (from defaults)
PREFECT_TASKS_DISABLE_CACHING='False' (from defaults)
PREFECT_TASKS_REFRESH_CACHE='False' (from defaults)
PREFECT_TASKS_RUNNER_PROCESS_POOL_MAX_WORKERS='None' (from defaults)
PREFECT_TASK_DEFAULT_RETRIES='0' (from defaults)
PREFECT_TASK_DEFAULT_RETRY_DELAY_SECONDS='0' (from defaults)
PREFECT_TASK_RUNNER_THREAD_POOL_MAX_WORKERS='None' (from defaults)
PREFECT_TASK_RUN_TAG_CONCURRENCY_SLOT_WAIT_SECONDS='10.0' (from defaults)
PREFECT_TASK_SCHEDULING_DEFAULT_STORAGE_BLOCK='None' (from defaults)
PREFECT_TASK_SCHEDULING_DELETE_FAILED_SUBMISSIONS='True' (from defaults)
PREFECT_TASK_SCHEDULING_MAX_RETRY_QUEUE_SIZE='100' (from defaults)
PREFECT_TASK_SCHEDULING_MAX_SCHEDULED_QUEUE_SIZE='1000' (from defaults)
PREFECT_TASK_SCHEDULING_PENDING_TASK_TIMEOUT='0:00:00' (from defaults)
PREFECT_TELEMETRY_ENABLE_RESOURCE_METRICS='True' (from defaults)
PREFECT_TELEMETRY_RESOURCE_METRICS_INTERVAL_SECONDS='10' (from defaults)
PREFECT_TEST_MODE='False' (from defaults)
PREFECT_TEST_SETTING='FOO' (from defaults)
PREFECT_UI_API_URL='<http://10.40.1.169:4200/api>' (from defaults)
PREFECT_UI_ENABLED='True' (from defaults)
PREFECT_UI_SERVE_BASE='/' (from defaults)
PREFECT_UI_STATIC_DIRECTORY='None' (from defaults)
PREFECT_UI_URL='<http://10.40.1.169:4200>' (from defaults)
PREFECT_UNIT_TEST_LOOP_DEBUG='True' (from defaults)
PREFECT_UNIT_TEST_MODE='False' (from defaults)
PREFECT_WORKER_CANCELLATION_POLL_SECONDS='120.0' (from defaults)
PREFECT_WORKER_DEBUG_MODE='False' (from defaults)
PREFECT_WORKER_ENABLE_CANCELLATION='False' (from defaults)
PREFECT_WORKER_HEARTBEAT_SECONDS='30.0' (from defaults)
PREFECT_WORKER_PREFETCH_SECONDS='10.0' (from defaults)
PREFECT_WORKER_QUERY_SECONDS='10.0' (from defaults)
PREFECT_WORKER_WEBSERVER_HOST='0.0.0.0' (from defaults)
PREFECT_WORKER_WEBSERVER_PORT='8080' (from defaults)Marvin
04/20/2026, 8:05 AMMarvin
04/20/2026, 8:08 AM# TaskRunRecorder — increase throughput and reduce latency
export PREFECT_SERVER_SERVICES_TASK_RUN_RECORDER_BATCH_SIZE=50
export PREFECT_SERVER_SERVICES_TASK_RUN_RECORDER_READ_BATCH_SIZE=10
export PREFECT_SERVER_SERVICES_TASK_RUN_RECORDER_FLUSH_INTERVAL=1
# EventPersister — align to higher throughput too
# Prefer canonical SERVER_ names; you can keep API_ as a fallback, but set these:
export PREFECT_SERVER_SERVICES_EVENT_PERSISTER_BATCH_SIZE=100
export PREFECT_SERVER_SERVICES_EVENT_PERSISTER_READ_BATCH_SIZE=10
export PREFECT_SERVER_SERVICES_EVENT_PERSISTER_FLUSH_INTERVAL=2
# Database pool — the server/services will open more concurrent DB connections under higher batching
# Prefer SERVER_ names over the deprecated top-level aliases you currently have set.
export PREFECT_SERVER_DATABASE_SQLALCHEMY_POOL_SIZE=20
export PREFECT_SERVER_DATABASE_SQLALCHEMY_MAX_OVERFLOW=20
export PREFECT_SERVER_DATABASE_SQLALCHEMY_POOL_TIMEOUT=45
# (optional but common)
export PREFECT_SERVER_DATABASE_SQLALCHEMY_POOL_RECYCLE=1800
If your actual peak concurrency is high (≈100+ concurrent tasks submitting state updates), use these instead:
export PREFECT_SERVER_SERVICES_TASK_RUN_RECORDER_BATCH_SIZE=200
export PREFECT_SERVER_SERVICES_TASK_RUN_RECORDER_READ_BATCH_SIZE=50
export PREFECT_SERVER_SERVICES_TASK_RUN_RECORDER_FLUSH_INTERVAL=0.5
export PREFECT_SERVER_SERVICES_EVENT_PERSISTER_BATCH_SIZE=250
export PREFECT_SERVER_SERVICES_EVENT_PERSISTER_READ_BATCH_SIZE=50
export PREFECT_SERVER_SERVICES_EVENT_PERSISTER_FLUSH_INTERVAL=1
export PREFECT_SERVER_DATABASE_SQLALCHEMY_POOL_SIZE=30
export PREFECT_SERVER_DATABASE_SQLALCHEMY_MAX_OVERFLOW=30
Notes on your current values
- Keeping TRR at batch_size=1 and flush_interval=5s can easily create multi-second to multi-minute lag when many tasks finish around the same time.
- Your EventPersister defaults (batch=20, read_batch=1, flush=5s) are okay but can backlog under load; bumping them helps the whole event pipeline.
- You currently rely on deprecated top-level DB pool aliases (PREFECT_SQLALCHEMY_POOL_SIZE, PREFECT_SQLALCHEMY_MAX_OVERFLOW). They still work, but I recommend switching to PREFECT_SERVER_DATABASE_SQLALCHEMY_* for clarity and future-proofing.
Operational tips
- After applying, restart the Prefect server services and watch logs for:
- Frequent “flushed N task runs” where N > 1, roughly every ~1s
- Absence of IntegrityError/OperationalError around TaskRunRecorder
- If workers are CPU-bound, consider also:
export PREFECT_FLOWS_HEARTBEAT_FREQUENCY=60
This reduces scheduling pressure on the heartbeat thread but is optional for the throughput fix.
If you can share your typical peak concurrency (total tasks that can complete per minute across the 4 workers), I can fine-tune these numbers further.Rupesh Raut
04/20/2026, 9:43 AMMarvin
04/20/2026, 9:43 AMMarvin
04/20/2026, 9:43 AMexport PREFECT_SERVER_SERVICES_TASK_RUN_RECORDER_BATCH_SIZE=50
export PREFECT_SERVER_SERVICES_TASK_RUN_RECORDER_READ_BATCH_SIZE=10
export PREFECT_SERVER_SERVICES_TASK_RUN_RECORDER_FLUSH_INTERVAL=1
export PREFECT_SERVER_SERVICES_EVENT_PERSISTER_BATCH_SIZE=100
export PREFECT_SERVER_SERVICES_EVENT_PERSISTER_READ_BATCH_SIZE=10
export PREFECT_SERVER_SERVICES_EVENT_PERSISTER_FLUSH_INTERVAL=2
- Observe for 15–30 minutes during load:
- Postgres: active connections, lock waits, statement durations
- Server logs: regular “flushed N task runs” with N > 1, reduced lag from task finish → state persisted
2) Phase 2: Only if you see “waited Xs for connection” or pool exhaustion
- Increase pool conservatively:
export PREFECT_SERVER_DATABASE_SQLALCHEMY_POOL_SIZE=15
export PREFECT_SERVER_DATABASE_SQLALCHEMY_MAX_OVERFLOW=15
export PREFECT_SERVER_DATABASE_SQLALCHEMY_POOL_TIMEOUT=45
- Recheck Postgres connection headroom before/after.
3) Phase 3: High concurrency (100+ tasks finishing per minute)
- Consider:
export PREFECT_SERVER_SERVICES_TASK_RUN_RECORDER_BATCH_SIZE=200
export PREFECT_SERVER_SERVICES_TASK_RUN_RECORDER_READ_BATCH_SIZE=50
export PREFECT_SERVER_SERVICES_TASK_RUN_RECORDER_FLUSH_INTERVAL=0.5
export PREFECT_SERVER_DATABASE_SQLALCHEMY_POOL_SIZE=20-30
export PREFECT_SERVER_DATABASE_SQLALCHEMY_MAX_OVERFLOW=20-30
- Tune up/down based on observed commit times and queue lag.
How to monitor safely
- Postgres (run on the DB)
- Active connections by app (helps attribute load):
SELECT application_name, count(*)
FROM pg_stat_activity
GROUP BY application_name ORDER BY 2 DESC;
- Long waits for connections/locks:
```
SELECT wait_event_type, wait_event, count(*)
FROM pg_stat_activity
WHERE state <> 'idle'Marvin
04/20/2026, 9:43 AM- Transaction volume trend:
SELECT date_trunc('minute', now()) AS ts,
(SELECT SUM(xact_commit+xact_rollback) FROM pg_stat_database) AS tx_total;
```
- Prefect server logs (DEBUG temporarily)
- Look for “flushed N task runs” frequency and any “IntegrityError”, “timeout waiting for connection”, or “dropping event”.
Risk mitigations
- Keep batch_size reasonable (50–200). Very large single commits can increase lock times.
- Ensure Postgres max_connections > sum of:
- Prefect server pool (and overflow) × number of server/service instances
- Other application pools
- 10–20 spare for admin/maintenance
- If you’re near max_connections or in a multi-replica topology, consider PgBouncer. In that case, we can size the Prefect pool smaller and let PgBouncer smooth spikes.
- Rollback is trivial: revert env vars and restart the server/services.
Bottom line
- The proposed batching changes primarily reduce transaction overhead while fixing your stuck RUNNING states; the net DB impact is usually equal or lower than before.
- Only increase DB pool if you observe contention. Start with batching, observe, then adjust pool sizes incrementally.
If you share:
- Number of server/service instances (pods/containers)
- Postgres max_connections
- Typical peak of tasks completing per minute
… I can give you a concrete ceiling for POOL_SIZE/MAX_OVERFLOW and confirm whether Phase 2 is needed.