Loam is pre-alpha: the engine core runs today; Live and Durable are in progress. See the roadmap

Blog/The Loam platform

Jobs on Loam: Celery, BullMQ, PySpark and Flink SQL

Our proposal for running the job frameworks teams already use on Loam with a configuration change, on one Rust queue core over TiKV, with durable steps from Resonate when you want them.

Jobs on Loam: Celery, BullMQ, PySpark and Flink SQL
On this page
  1. One core: operon-jobs
  2. Where the state lives
  3. Semantics, stated plainly
  4. Celery
  5. BullMQ v6
  6. Durable mode, with Resonate's decorators
  7. PySpark on Sail
  8. Flink SQL on RisingWave
  9. Where this stands

Proposal This post describes a design under review. No jobs code exists yet.

Every AI application grows background work: embedding backfills, document parsing, email, nightly feature builds, streaming aggregations. Teams already have code for it, written against Celery, BullMQ, PySpark or Flink, and a broker to run it: Redis, RabbitMQ, a Spark cluster, a Flink cluster. We do not want them to rewrite that code. The proposal is to run each family on Loam with as little change as the family allows, on one Rust core, and to make durability opt-in with Resonate's own decorators.

WorkloadHow it would run on LoamChange to your code
CeleryA kombu transport, a result backend and a beat scheduler, in a Python packageThe broker URL and result backend become loam://…
BullMQ v6BullMQ's own classes, bound to a Loam implementation of BullMQ's pluggable backendThe import changes from bullmq to @loam/bullmq
PySparkSail, run by Loam as a Spark Connect server per namespaceSparkSession.builder.remote("sc://…")
Flink SQLRisingWave, Loam's companion stream processorA dialect port. DataStream jobs run unmodified on Apache Flink

One core: operon-jobs

Under all four sits one Rust crate with one protobuf surface, loam.jobs.v1, served over connect-rust like the Live sync API. Python and TypeScript clients are generated from it, and the Celery and BullMQ adapters are thin layers over those clients.

#[async_trait]
pub trait Jobs: Send + Sync {
    async fn enqueue(&self, cx: &Ctx, queue: &QueueId, task: TaskSpec, opts: EnqueueOpts)
        -> Result<Enqueued, JobsError>;
    async fn lease(&self, cx: &Ctx, req: LeaseRequest) -> Result<Vec<Lease>, JobsError>;
    async fn extend(&self, cx: &Ctx, token: &LeaseToken, by: Duration) -> Result<Deadline, JobsError>;
    async fn complete(&self, cx: &Ctx, token: &LeaseToken, outcome: Outcome)
        -> Result<Completed, JobsError>;
    async fn schedule(&self, cx: &Ctx, spec: ScheduleSpec) -> Result<ScheduleId, JobsError>;
    async fn flow(&self, cx: &Ctx, graph: FlowSpec) -> Result<FlowHandle, JobsError>;
    async fn submit_engine_job(&self, cx: &Ctx, job: EngineJob, opts: SubmitOpts)
        -> Result<RunId, JobsError>;
    fn watch(&self, cx: &Ctx, what: Selector, from: Cursor)
        -> BoxStream<'static, Result<Event, JobsError>>;
    // … plus bulk enqueue, progress reports, engine control and admin calls
}

The namespace always comes from the caller's credential (Ctx), never from a request body. Every write that depends on holding a job carries a LeaseToken with a fencing epoch.

Where the state lives

StateStore
Job records, ready and delayed indexes, leases, dedup keys, rate buckets, chord countersA TiKV keyspace. A local redb file in operon dev only; every other deployment needs TiKV
Payloads and results over 16 KiBThe object store
The queue's event logA Loam stream per queue, written through an outbox with an idempotent producer
Schedules, flows and durable-mode tasksThe embedded Resonate server

Every queue operation is a small multi-key transaction: dequeue means read the ready index, lock the job and write its lease. TiKV's transactions give exactly that. No Redis, RabbitMQ or ZooKeeper is involved.

Semantics, stated plainly

  • Delivery is at least once. A job goes to one worker at a time. It is delivered again if its lease expires (a crashed or stalled worker) or the worker gives it back.
  • The outcome is recorded exactly once. Completion is fenced by the lease epoch. If a stalled worker and its replacement both finish a job, only the holder of the current epoch can record the result; the other gets Fenced. A result, a retry count or a parent's child counter is never applied twice.
  • Effectively-once needs idempotent tasks or durable mode. A task's own side effects happen at least once in queue mode, exactly as they do on Redis or RabbitMQ today. Loam never claims exactly-once execution.

Celery

Celery has three plug points, and the package implements all three. The transport maps kombu calls onto the queue: _put is enqueue, _get is a long-polled lease, basic_ack is complete, basic_reject is release or dead-letter. The lease is a server-side visibility timeout, and the transport extends it every third of the timeout while a task runs, so a long acks_late task on a live worker is never redelivered, which removes Redis's "a task longer than the visibility timeout runs twice" problem. ETA tasks go into the delayed index instead of sitting unacked in a worker's memory. The result backend implements Celery's key-value primitives including incr, so chords use Celery's native counter in one TiKV transaction instead of the polling chord_unlock task. Beat upserts every schedule entry as a server-side schedule, so running beat twice cannot duplicate a tick.

# celeryconfig.py: the only change for a queue-mode app
broker_transport = "loam_celery.transport:Transport"
broker_url = "loam://jobs.example:7720/my-namespace"
result_backend = "loam://jobs.example:7720/my-namespace"
beat_scheduler = "loam"

Some Celery features are switched on by a hard-coded list of broker types rather than by capabilities. Remote control and events work; mingle and gossip stay off.

BullMQ v6

We expected to have to emulate Redis for BullMQ. Reading the source changed the plan: BullMQ v6 has a pluggable backend interface, IQueueBackend, about 81 methods that describe every operation the Queue, Worker and Job classes need, with a Redis and a PostgreSQL implementation. So @loam/bullmq is a backend factory for the unmodified bullmq package, plus a module that re-exports BullMQ's classes already bound to it. The code users run is BullMQ's own, so the API cannot drift. Emulating Redis instead would mean matching 49 Lua scripts and 49 Redis commands, forever, and inviting every other Redis workload.

The gate is BullMQ's own backend-neutral test suite passing against Loam. It supports v6 only (v5 has no backend interface), and not BullMQ Pro's commercial features.

Durable mode, with Resonate's decorators

Queue mode runs tasks as the framework always has. Durable mode is opt-in per task, using Resonate's own SDKs, not a new decorator system. The framework keeps admission (queue, priority, rate limit, delay), and the task body becomes a Resonate function whose promise id is derived from the job id. A redelivered job then resumes the same execution: finished steps return their recorded results.

@resonate.register(name="charge", version=1)
async def charge(ctx, order_id: str):
    hold = await ctx.run(reserve, order_id)
    receipt = await ctx.run(capture, hold)      # checkpointed: skipped on resume once recorded; capture must still be idempotent
    await ctx.run(send_receipt, receipt)
    return receipt

@app.task(bind=True, acks_late=True)            # still a Celery task
def charge_task(self, order_id):
    handle = resonate.run(f"celery:{self.request.id}", charge, order_id)
    return asyncio.run(handle.result())

Queue-mode jobs create no promises, which keeps millions of short jobs out of the durable store. Schedules and flows (Celery's chain, group and chord, BullMQ's parent and child flows) run on Resonate either way: a flow is a durable function whose joins are promises, so a crash mid-flow resumes from its last finished step.

PySpark on Sail

Sail is a Spark Connect server written in Rust on DataFusion. PySpark clients connect to it unchanged. Loam would run one Sail server per namespace (Sail runs Python UDFs in its own process, so sharing one across tenants would run one tenant's Python beside another's data), start it on the first connection, scale it to zero when idle, and route sc:// connections to it through an authenticating proxy. Tables live on the object store as Iceberg through Lakekeeper.

Sail does not cover everything. Spark Connect does not carry RDDs or SparkContext at all. Sail's own CI reports about 91% of Spark 3.5's Connect tests passing and about 71% of Spark 4.2's. Structured Streaming is not ready. So there is a fallback: jobs Sail cannot run go to Apache Spark on Kubernetes, unmodified, reading the same Iceberg tables. A loam spark check command would run a job on Sail and report which calls it could not handle, before a team moves anything. Resonate only wraps whole Spark actions as steps; it never enters Spark's own task retries.

For streaming SQL, the proposal moves Flink SQL jobs to RisingWave, which Loam already plans as its companion stream processor: Rust, Apache-2.0, state in object storage, Iceberg sources and sinks. RisingWave speaks Postgres-dialect SQL, not Flink SQL, so "little change" means a porting table with a worked example per construct, not a translator. DataStream jobs have no Rust replacement; they run on Apache Flink with the Flink Kubernetes Operator, checkpointing to the object store, and Loam manages only their lifecycle (deploy, savepoint, upgrade, rollback) as durable workflows. Flink's checkpoints stay the job's only durable state. The Sail and RisingWave deep dive covers both projects.

Where this stands

PartStatus
loam.jobs.v1, the queue core, the Celery transport, result backend and beat, TiKV job storePlanned
@loam/bullmq on BullMQ v6's backend interface, flows on Resonate promisesPlanned
PySpark on Sail per namespace, the Spark Connect proxy, the Spark fallbackPlanned
Flink SQL on RisingWave, DataStream lifecycle on the Flink operatorPlanned

The first phase is Celery, because it exercises the whole core: leases, fencing, delays, retries, dead letters, schedules on Resonate and the TiKV store. Its gate is Celery's own integration suite passing against operon dev. Jobs depend on the embedded durable server and the TiKV client layer, both of which exist. The next post covers the other databases teams bring: Postgres and MySQL.

More from the blog