CarbonQueue: Carbon-Aware Batch Job Scheduling System¶
Rajesh Nandipaty · December 7, 2025 · CarbonQueue schedules deferrable jobs in the cleanest feasible window before their deadline.
CarbonQueue can produce different CO2e emissions depending on when it runs, as the grid’s generation mix varies throughout the day. A batch job due tomorrow morning can often be shifted away from higher-intensity periods to cleaner overnight or midday hours. CarbonQueue forecasts hourly grid carbon intensity, evaluates each feasible start time before the deadline, and selects the window with the lowest mean forecast intensity. It then reports the estimated emissions avoided relative to immediate execution.
This notebook runs the full system on a 28-day synthetic intensity trace with known structure. The forecaster and scheduler operate normally, while the known realized trace enables evaluation against a perfect-foresight oracle rather than relying on asserted savings. In live mode, the synthetic trace is replaced by the Electricity Maps API; all other components remain unchanged.
1. Configuration¶
The grid zone, job deadline, and forecast horizon are configured here. In simulation, these values come from the bundled synthetic trace. In live mode, the grid zone determines the Electricity Maps client configuration.
import os, sys, json, time
sys.path.insert(0, os.path.abspath("."))
import numpy as np, pandas as pd, matplotlib.pyplot as plt
from IPython.display import Image, display
plt.rcParams.update({"figure.dpi": 120, "axes.grid": True, "grid.alpha": 0.25})
os.makedirs("nbfigs", exist_ok=True) # Scratch dir for the notebook's inline figures.
from carbonqueue import intensity, forecast as fc, scheduler, sim, report
from carbonqueue.store import JobStore
ZONE, DEADLINE_H, HORIZON, SEED = "US-NY-NYIS", 24.0, 48, 7
OUT = "outputs"; os.makedirs(OUT, exist_ok=True)
print(f"Zone: {ZONE} · Deadline: {DEADLINE_H:.0f} h · Forecast horizon: {HORIZON} h · Seed: {SEED}")
Zone: US-NY-NYIS · Deadline: 24 h · Forecast horizon: 48 h · Seed: 7
2. Carbon Signal¶
The input is an hourly series of grid carbon intensity measured in gCO2e/kWh. The synthetic generator captures typical grid patterns, including an overnight plateau, a midday dip as solar generation increases, an evening ramp as solar generation declines and demand peaks, weekday and weekend variation, and autocorrelated noise. It serves as a substitute for a cached historical export. The plot shows one week of the trace. The variation between daily troughs and peaks represents the scheduling opportunity the system seeks to exploit.
trace = intensity.synthetic_trace(days=28, seed=SEED)
wk = trace.iloc[:24 * 7]
fig, ax = plt.subplots(figsize=(11, 3.2))
ax.plot(wk.index, wk.values, lw=1.1, color="#2a9d5c")
ax.set_ylabel("gCO2e / kWh"); ax.set_title("Grid Carbon Intensity. One week of the trace.")
fig.savefig("nbfigs/trace_week.png", bbox_inches="tight"); plt.close(fig)
display(Image("nbfigs/trace_week.png"))
print(f"Trace: {len(trace)} hours, {trace.min():.0f}-{trace.max():.0f} gCO2e/kWh. "
f"Daily swing: ~{trace.groupby(trace.index.hour).mean().max() - trace.groupby(trace.index.hour).mean().min():.0f}")
Trace: 672 hours, 107-404 gCO2e/kWh. Daily swing: ~185
3. Durable Queue and Event Log¶
Jobs are stored in SQLite, with each state change appended to an event log. This preserves job state across restarts and provides an auditable record of scheduling decisions. After three failed attempts, a job transitions to a dead-letter state rather than being retried indefinitely. These components implement task-queue and dead-letter patterns described in Designing Data-Intensive Applications at single-machine scale. The following demonstration submits and schedules a job, then deliberately fails it three times to illustrate the dead-letter transition and resulting event trail.
store = JobStore(":memory:")
jid = store.submit("nightly_retrain", "python train.py", est_hours=3,
deadline_ts=time.time() + DEADLINE_H * 3600, device_watts=350)
store.set_schedule(jid, time.time() + 6 * 3600, mean_intensity=210.0, now_intensity=330.0)
for _ in range(3):
store.fail(jid, "simulated worker crash") # third failure -> dead-letter
events = pd.DataFrame(store.all_jobs())[["name", "state", "attempts"]]
log = pd.read_sql_query("SELECT seq, kind, job_id FROM events ORDER BY seq", store.conn)
print(events.to_string(index=False)); print("\nEvent log:")
print(log.to_string(index=False))
name state attempts nightly_retrain dead 3 Event log: seq kind job_id 1 submitted 6a119de97bb3 2 scheduled 6a119de97bb3 3 failed 6a119de97bb3 4 failed 6a119de97bb3 5 dead_lettered 6a119de97bb3
4. Forecasting¶
The scheduler uses forecasts to rank future execution windows. Two models are evaluated using the same procedure: a seasonal-naive baseline that repeats the value from 24 hours earlier, and a ridge regression using lag values at 1, 2, 3, 24, 48, and 168 hours together with calendar features. Model selection uses walk-forward validation with MAE as the evaluation metric, following the time-series methodology described in Hands-On Machine Learning, Chapter 15. The ridge model is selected only when it improves on the seasonal-naive baseline. Otherwise, the baseline provides the fallback model. The plot compares a day-ahead forecast with the corresponding observed values.
acc = fc.walk_forward_mae(trace.iloc[:14 * 24 + 24 * 7], train_days=14)
print("Walk-forward day-ahead MAE:", acc)
cut = 20 * 24
hist, actual = trace.iloc[:cut], trace.iloc[cut:cut + 24]
ridge = fc.RidgeForecaster().fit(hist).predict(hist, 24)
naive = fc.seasonal_naive(hist, 24)
fig, ax = plt.subplots(figsize=(10, 3.2))
ax.plot(actual.index, actual.values, color="#333", lw=2, label="actual")
ax.plot(ridge.index, ridge.values, color="#2a9d5c", lw=1.6, label="ridge")
ax.plot(naive.index, naive.values, color="#c0392b", lw=1.2, ls="--", label="seasonal naive")
ax.set_ylabel("gCO2e / kWh"); ax.set_title("Day-Ahead Forecast vs Actual"); ax.legend(fontsize=8)
fig.savefig("nbfigs/forecast.png", bbox_inches="tight"); plt.close(fig)
display(Image("nbfigs/forecast.png"))
Walk-forward day-ahead MAE: {'seasonal_naive_mae': 29.3, 'ridge_mae': 26.7, 'windows_evaluated': 6}
5. Window Selection¶
Given a forecast, the scheduler evaluates every feasible start hour between the current time and the latest start time that satisfies the deadline. For each candidate, it computes the mean forecast carbon intensity over the job duration and selects the window with the lowest value. The plot shows a three-hour job submitted during an evening peak and compares the mean forecast intensity across candidate start times. The scheduler selects the lowest-intensity window rather than the submission hour.
at = 20 * 24 + 18 # A submission at hour 18 (evening peak).
dur = 3
hist = trace.iloc[:at]
fcast = fc.RidgeForecaster().fit(hist).predict(hist, HORIZON)
now_ts = trace.index[at].timestamp()
deadline_ts = now_ts + DEADLINE_H * 3600
w = scheduler.choose_window(fcast, now_ts, dur, deadline_ts)
latest = pd.Timestamp(deadline_ts, unit="s") - pd.Timedelta(hours=dur)
cands = [t for t in fcast.index if trace.index[at] <= t <= latest]
means = [float(fcast.loc[t:t + pd.Timedelta(hours=dur - 1)].mean()) for t in cands]
fig, ax = plt.subplots(figsize=(10, 3.2))
ax.plot(cands, means, color="#2a9d5c", marker="o", ms=3)
ax.axvline(pd.Timestamp(w["start"]), color="#2a9d5c", ls="--", label=f"chosen: {pd.Timestamp(w['start']):%a %H:%M}")
ax.axvline(trace.index[at], color="#c0392b", ls=":", label=f"submitted: {trace.index[at]:%a %H:%M}")
ax.set_ylabel("mean forecast gCO2e/kWh over 3 h"); ax.set_title("Candidate Start Windows for One Job"); ax.legend(fontsize=8)
fig.savefig("nbfigs/window_scan.png", bbox_inches="tight"); plt.close(fig)
display(Image("nbfigs/window_scan.png"))
print(f"Submitted at {trace.index[at]:%A %H:%M} and scheduled for {pd.Timestamp(w['start']):%A %H:%M}. "
f"Forecast mean of {w['mean_intensity']:.0f} vs {w['now_intensity']:.0f} now.")
Submitted at Sunday 18:00 and scheduled for Monday 12:00. Forecast mean of 202 vs 303 now.
6. Batch Simulation¶
Twelve jobs drawn from six workload profiles are submitted at random times over the 28-day simulation, each with a 24-hour deadline. At each submission, the forecaster is trained using only the history available at that time, so the scheduler has no access to future observations. Final emissions are calculated from the realized carbon-intensity trace rather than the forecast.
placements = sim.simulate_batch(trace, jobs=12, deadline_hours=DEADLINE_H, seed=SEED, train_days=14)
summary = sim.summarize(trace, placements, DEADLINE_H, f"synthetic(28d, seed {SEED})", acc)
df = pd.DataFrame(placements)
df[["name", "submitted", "start", "delay_hours", "forecast_mean", "grams_if_run_now", "grams_actual", "grams_avoided"]]
| name | submitted | start | delay_hours | forecast_mean | grams_if_run_now | grams_actual | grams_avoided | |
|---|---|---|---|---|---|---|---|---|
| 0 | nightly_model_retrain_00 | 2025-10-20 16:00:00 | 2025-10-21 12:00:00 | 20.0 | 226.3 | 323.5 | 173.1 | 150.4 |
| 1 | full_db_backup_01 | 2025-10-22 17:00:00 | 2025-10-23 12:00:00 | 19.0 | 216.4 | 75.7 | 52.6 | 23.1 |
| 2 | ci_regression_suite_02 | 2025-10-23 12:00:00 | 2025-10-23 14:00:00 | 2.0 | 215.9 | 20.0 | 17.9 | 2.1 |
| 3 | media_transcode_batch_03 | 2025-10-23 16:00:00 | 2025-10-24 12:00:00 | 20.0 | 214.3 | 198.1 | 114.0 | 84.1 |
| 4 | embedding_reindex_04 | 2025-10-26 23:00:00 | 2025-10-27 12:00:00 | 13.0 | 195.1 | 152.0 | 119.7 | 32.3 |
| 5 | report_render_farm_05 | 2025-10-27 10:00:00 | 2025-10-27 13:00:00 | 3.0 | 182.2 | 12.4 | 12.6 | -0.2 |
| 6 | nightly_model_retrain_06 | 2025-10-28 04:00:00 | 2025-10-28 11:00:00 | 7.0 | 223.8 | 310.5 | 154.6 | 155.9 |
| 7 | full_db_backup_07 | 2025-10-29 08:00:00 | 2025-10-29 13:00:00 | 5.0 | 183.0 | 68.7 | 46.1 | 22.6 |
| 8 | ci_regression_suite_08 | 2025-10-30 02:00:00 | 2025-10-30 13:00:00 | 11.0 | 183.9 | 22.5 | 17.4 | 5.0 |
| 9 | media_transcode_batch_09 | 2025-10-30 18:00:00 | 2025-10-31 12:00:00 | 18.0 | 204.5 | 232.6 | 108.8 | 123.8 |
| 10 | embedding_reindex_10 | 2025-10-31 05:00:00 | 2025-10-31 13:00:00 | 8.0 | 201.0 | 180.4 | 101.2 | 79.2 |
| 11 | report_render_farm_11 | 2025-11-01 07:00:00 | 2025-11-01 12:00:00 | 5.0 | 171.7 | 20.2 | 11.8 | 8.4 |
7. Results¶
Across the batch, CarbonQueue reduces modeled emissions with a modest median execution delay by shifting jobs away from higher-intensity evening periods toward lower-intensity overnight and midday windows. The trace figure illustrates this behavior, with submission times concentrated near higher-intensity periods and execution windows occurring during lower-intensity periods. The per-job figure shows the contribution of each job to the overall emissions reduction.
report.fig_trace(trace.iloc[14 * 24 - 24:], placements, f"{OUT}/carbonqueue-trace.png")
report.fig_savings(placements, f"{OUT}/carbonqueue-savings.png")
print(f"Avoided {summary['pct_avoided']:.1f}% of modeled emissions "
f"({summary['total_g_scheduled']:.0f} g scheduled vs {summary['total_g_run_now']:.0f} g run-now) "
f"at a median delay of {summary['median_delay_hours']:.1f} h.")
from IPython.display import Image, display
display(Image(f"{OUT}/carbonqueue-trace.png")); display(Image(f"{OUT}/carbonqueue-savings.png"))
Avoided 42.5% of modeled emissions (930 g scheduled vs 1617 g run-now) at a median delay of 9.5 h.
8. Evaluation Against an Oracle¶
Comparison with immediate execution demonstrates the benefit of carbon-aware scheduling, but does not establish whether the forecast is sufficiently accurate. To evaluate forecast performance, CarbonQueue is compared with a perfect-foresight oracle that selects the lowest-intensity feasible window from the realized trace. The capture ratio measures the proportion of the oracle's achievable emissions reduction captured by the forecast-based scheduler. A deadline-slack sweep shows that emissions reductions increase with greater scheduling flexibility, while the gap between forecast-based and oracle performance also widens. This reflects the increasing difficulty of identifying the lowest-intensity window as the scheduling horizon expands.
o = summary["oracle"]
print(f"Run-now avoidable ceiling (oracle): {o['pct_avoided_oracle']:.1f}%")
print(f"CarbonQueue (forecast): {summary['pct_avoided']:.1f}%")
print(f"Capture ratio: {o['capture_ratio']:.3f} ")
sweep = sim.slack_sweep(trace, seed=SEED, train_days=14, jobs=12)
report.fig_slack_sweep(sweep, f"{OUT}/carbonqueue-slack-sweep.png")
display(sweep)
display(Image(f"{OUT}/carbonqueue-slack-sweep.png"))
Run-now avoidable ceiling (oracle): 44.4% CarbonQueue (forecast): 42.5% Capture ratio: 0.956
| deadline_hours | pct_avoided | pct_avoided_oracle | median_delay_hours | |
|---|---|---|---|---|
| 0 | 6 | 7.9 | 8.0 | 2.0 |
| 1 | 12 | 21.5 | 24.7 | 5.0 |
| 2 | 24 | 42.5 | 44.4 | 9.5 |
| 3 | 48 | 42.6 | 49.6 | 12.0 |
| 4 | 72 | 42.6 | 51.2 | 12.0 |
9. Summary¶
A machine-readable summary of the results is written to outputs/carbonqueue-summary.json.
report.write_summary(f"{OUT}/carbonqueue-summary.json", summary)
print(json.dumps({k: v for k, v in summary.items() if k != "jobs"}, indent=2, default=str))
{
"mode": "simulation",
"trace_source": "synthetic(28d, seed 7)",
"trace_hours": 672,
"deadline_hours": 24.0,
"jobs_placed": 12,
"forecast_eval": {
"seasonal_naive_mae": 29.3,
"ridge_mae": 26.7,
"windows_evaluated": 6
},
"median_delay_hours": 9.5,
"total_g_run_now": 1616.6,
"total_g_scheduled": 929.8,
"total_g_avoided": 686.8,
"pct_avoided": 42.5,
"oracle": {
"total_g_oracle": 898.1,
"pct_avoided_oracle": 44.4,
"capture_ratio": 0.956
}
}
10. Limitations¶
- Energy consumption is modeled rather than measured. Job energy is estimated as duration multiplied by nameplate power, which does not account for utilization or idle consumption. RAPL counters or Kepler metrics on Kubernetes could replace this estimate with measured energy, allowing avoided emissions to be calculated from observed energy use.
- The intensity signal represents grid-average rather than marginal carbon intensity. These measures can differ in which hours they identify as lowest-carbon. Marginal intensity data is available through sources with different access requirements.
- Emissions reductions depend on the characteristics of the trace. The approximately 42% reduction observed in this simulation reflects a trace with substantial daily variation and relatively generous deadlines. A flatter intensity profile or tighter deadlines would produce smaller reductions, as demonstrated by the slack sweep. The broader result is that the scheduler shifts work toward lower-intensity periods within the available scheduling window.
- The optimization objective is limited to carbon. CarbonQueue minimizes carbon intensity and reports the resulting emissions reduction. Electricity prices may peak at different times, so a joint cost and carbon optimization objective remains outside the scope of this implementation.