Mojo Concurrency

Concurrency means structuring a program to handle multiple tasks that can make progress at overlapping times. Parallelism (covered earlier) runs tasks simultaneously on multiple CPU cores. Concurrency is the broader concept — tasks take turns sharing resources, which matters most for I/O-bound work like network calls, file reads, and database queries where the CPU spends most of its time waiting.

Concurrency vs Parallelism

Concurrency (one worker, multiple tasks):
  Task A: ████░░░░████░░░░████
  Task B: ░░░░████░░░░████░░░░
  Time:   ────────────────────→
  One worker switches between tasks. Neither blocks the other.

Parallelism (multiple workers):
  Task A: ████████████████████  ← CPU core 0
  Task B: ████████████████████  ← CPU core 1
  Time:   ────────────────────→
  Two workers run simultaneously.

Both use the wall-clock time efficiently.
The difference is HOW idle time is used.

When Each Applies

Use PARALLELISM for:          Use CONCURRENCY for:
  CPU-bound work                I/O-bound work
  ─────────────────             ────────────────────────
  Numerical computation         Waiting for network responses
  Image processing              Reading files from disk
  Matrix multiply               Waiting for database queries
  SIMD operations               Handling many client connections

  The CPU is always busy.       The CPU is mostly waiting.

Python async/await via Interop

Mojo does not yet have its own async runtime. For I/O-bound concurrency today, use Python's asyncio library through Mojo's Python interop layer.

from python import Python

fn main() raises:
    var asyncio = Python.import_module("asyncio")

    # Define an async function in Python syntax via exec
    Python.evaluate("""
import asyncio
import time

async def fetch_data(name, delay):
    print(f"  Starting: {name}")
    await asyncio.sleep(delay)   # simulates waiting for I/O
    print(f"  Done:     {name}")
    return f"result from {name}"

async def main_async():
    tasks = [
        fetch_data("Task A", 1.0),
        fetch_data("Task B", 0.5),
        fetch_data("Task C", 1.5),
    ]
    results = await asyncio.gather(*tasks)
    for r in results:
        print(r)

asyncio.run(main_async())
""")

Output (tasks interleave, total time ≈ 1.5 s instead of 3 s sequential):

  Starting: Task A
  Starting: Task B
  Starting: Task C
  Done:     Task B   ← finishes first (shortest wait)
  Done:     Task A
  Done:     Task C   ← finishes last (longest wait)
result from Task A
result from Task B
result from Task C
Timeline:
  t=0.0: Task A, B, C all start (concurrent)
  t=0.5: Task B finishes (0.5s wait done)
  t=1.0: Task A finishes (1.0s wait done)
  t=1.5: Task C finishes (1.5s wait done)
  Total: 1.5s  (vs 3.0s sequential — 2× faster)

Thread-Based Concurrency via Python threading

from python import Python

fn main() raises:
    Python.evaluate("""
import threading
import time

results = []
lock = threading.Lock()

def worker(name, seconds):
    time.sleep(seconds)
    with lock:
        results.append(f"{name} completed after {seconds}s")

threads = [
    threading.Thread(target=worker, args=("Worker 1", 1)),
    threading.Thread(target=worker, args=("Worker 2", 2)),
    threading.Thread(target=worker, args=("Worker 3", 1)),
]

for t in threads:
    t.start()

for t in threads:
    t.join()

for r in results:
    print(r)
""")

Mojo's Built-In Concurrency: parallelize

For CPU-bound concurrent tasks, Mojo's parallelize is the native tool. It distributes work across cores with no Python overhead.

from algorithm import parallelize
from memory import UnsafePointer

fn main():
    let n = 8
    var results = UnsafePointer[Int].alloc(n)
    for i in range(n):
        results.init_pointee_copy(0)

    @parameter
    fn heavy_task(i: Int):
        # Simulate a computation-heavy task per slot
        var acc = 0
        for j in range(1_000_000):
            acc += j * i
        results[i] = acc % 1000

    parallelize[heavy_task](n)   # runs all 8 tasks concurrently across cores

    for i in range(n):
        print("Task", i, "result:", results[i])

    for i in range(n):
        (results + i).destroy_pointee()
    results.free()

Synchronization: Avoiding Shared State Problems

The core concurrency hazard — shared mutable state:

  Task A reads shared_count (= 5)
  Task B reads shared_count (= 5)
  Task A writes shared_count = 5 + 1 = 6
  Task B writes shared_count = 5 + 1 = 6   ← lost Task A's update!

  Expected: 7   Actual: 6   → data race

Solutions:
  1. No sharing: give each task its own independent data (best)
  2. Atomic operations: reads and writes as one uninterruptible step
  3. Locks / mutexes: only one task accesses shared data at a time

Mojo's Approach: Safety Through Ownership

Mojo's ownership system prevents data races by design:
  - Only one mutable reference (inout) exists at any time
  - Multiple immutable references (borrowed) are fine simultaneously
  - parallelize() gives each iteration its own exclusive index
    → no two iterations share a write target

  This is safer than Python's threading (GIL is loose,
  not a true race prevention mechanism) and safer than
  raw C threads (no compiler enforcement).

Key Takeaways

Concurrency structures a program to handle multiple tasks at overlapping times; parallelism runs them simultaneously on multiple cores. Concurrency shines for I/O-bound work; parallelism for CPU-bound computation. Use Python's asyncio through Mojo's interop for network and file I/O concurrency today. Use Mojo's parallelize for CPU-bound concurrent work natively. The root problem in concurrency is shared mutable state — avoid it by giving each task independent data. Mojo's ownership system enforces this at compile time, making many concurrency bugs impossible to write.

Leave a Comment

Your email address will not be published. Required fields are marked *