惯性聚合 高效追踪和阅读你感兴趣的博客、新闻、科技资讯
阅读原文 在惯性聚合中打开

推荐订阅源

Stack Overflow Blog
Stack Overflow Blog
N
News | PayPal Newsroom
阮一峰的网络日志
阮一峰的网络日志
月光博客
月光博客
T
Tailwind CSS Blog
博客园 - 叶小钗
博客园 - 【当耐特】
Apple Machine Learning Research
Apple Machine Learning Research
B
Blog RSS Feed
Know Your Adversary
Know Your Adversary
P
Privacy International News Feed
cs.CL updates on arXiv.org
cs.CL updates on arXiv.org
Project Zero
Project Zero
美团技术团队
雷峰网
雷峰网
Martin Fowler
Martin Fowler
P
Privacy & Cybersecurity Law Blog
T
The Blog of Author Tim Ferriss
S
Schneier on Security
V
V2EX
Cisco Talos Blog
Cisco Talos Blog
Blog — PlanetScale
Blog — PlanetScale
G
GRAHAM CLULEY
J
Java Code Geeks
V
Visual Studio Blog
N
Netflix TechBlog - Medium
Threat Intelligence Blog | Flashpoint
Threat Intelligence Blog | Flashpoint
Cyberwarzone
Cyberwarzone
Recent Announcements
Recent Announcements
C
CXSECURITY Database RSS Feed - CXSecurity.com
让小产品的独立变现更简单 - ezindie.com
让小产品的独立变现更简单 - ezindie.com
N
News and Events Feed by Topic
Forbes - Security
Forbes - Security
GbyAI
GbyAI
奇客Solidot–传递最新科技情报
奇客Solidot–传递最新科技情报
The Hacker News
The Hacker News
Application and Cybersecurity Blog
Application and Cybersecurity Blog
The Cloudflare Blog
腾讯CDC
爱范儿
爱范儿
Exploit-DB.com RSS Feed
Exploit-DB.com RSS Feed
OSCHINA 社区最新新闻
OSCHINA 社区最新新闻
罗磊的独立博客
T
Threatpost
大猫的无限游戏
大猫的无限游戏
NISL@THU
NISL@THU
Cloudbric
Cloudbric
C
CERT Recently Published Vulnerability Notes
C
Check Point Blog
Microsoft Azure Blog
Microsoft Azure Blog

DEV Community

Authentication Security Deep Dive: From Brute Force to Salted Hashing (With Java Examples) Why AI Systems Don’t Fail — They Drift Spilling beans for how i learn for exam😁"Reinforcement Learning Cheat Sheet" I Replaced Chrome with Safari for AI Browser Automation. Here's What Broke (and What Finally Worked) How Python Borrows Other People's Work The $40 Architecture: Processing 1 Billion API Requests with 99.99% Uptime Vibe Coding: A Workflow Guide (From Zero to SaaS) Most webhook security guides protect the wrong side. The scary part is delivery. Headless CMS for TanStack Start: Build a Blog with Cosmic EU Age Verification App "Hacked in 2 Minutes" — What Actually Happened Comfy Cloud’s delete function does not actually remove files Running AI Models on GPU Cloud Servers: A Beginner Guide Event-driven media intelligence with AWS Step Functions and Bedrock I scored 500 AI prompts across 8 quality dimensions — here's what broke How to Call Google Gemini API from Next.js (Free Tier, No Backend Needed) The Portal Protocol: Reclaiming Human Connection in the Age of AI How to Fix Your Team's Scattered Knowledge Problem With a Self-Hosted Forum Intro to tc Cloud Functors: A Graph-First Mental Model for the Modern Cloud Designing Multi-Tenant Backends With Both Ownership and Team Access I Built a Neumorphic CSS Library with 77+ Components — Here's What I Learned PostgreSQL Performance Optimization: Why Connection Pooling Is Critical at Scale Cómo construí un SaaS multi-rubro para gestionar expensas en Argentina con FastAPI + Vue 3 🚀 I Built an Ethical Hacking Scanner Tool – Open Source Project I Replaced /usage and /context in Claude Code With a Single Statusline A Pythonic Way to Handle Emails (IMAP/SMTP) with Auto-Discovery and AI-Ready Design I Collected 8.9 Million Polymarket Price Points — Here's What I Found About How Markets Really Move EcoTrack AI — Carbon Footprint Tracker & Dashboard Everyone's Using AI. No One Agrees How. 5 self-hosted ebook managers worth trying in 2026 Building Your First AI Agent with LangChain: From Chatbot to Autonomous Assistant Common SOC 2 Failures (Real World) Stop Vibe-Checking Your AI App: A Practical Guide to Evals How to Use SonarQube and SonarScanner Locally to Level Up Your Code Quality Your Next To-Do App Is Dead — I Replaced Mine with an OpenClaw AI Sign a Nostr event in 60 lines of Python using coincurve — no nostr-sdk, no nbxplorer, no rust toolchain ITGC Audit Explained Like You’re in Big 4 Patch Tuesday abril 2026: Microsoft parcha 163 vulnerabilidades y un zero-day en SharePoint Stop scraping everything: a better way to track competitor price changes Listing on MCPize + the Official MCP Registry while routing payments OUTSIDE the marketplace — how I kept 100% of my x402 revenue Building an AI-Powered Risk Intelligence System Using Serverless Architecture Why We Ripped Function Overloading Out of Our AI Toolchain Testing AI-Generated Code: How to Actually Know If It Works SaaS Churn Is Killing Your Business. Here Is What to Do About It (Without a Support Team) The Speed of AI Is No Longer Linear - And Self-Improving Models Are Why How to Implement RBAC for MCP Tools: A Practical Guide for Engineering Teams From Standard Quote to Persuasive Proposal: AI Automation for Arborists I built a CLI that scaffolds complete multi-tenant SaaS apps Axios CVE-2025–62718: The Silent SSRF Bug That Could Be Hiding in Your Node.js App Right Now The dashboard that ended our friendship Data Pipelines Explained Simply (and How to Build Them with Python) The Hidden Cost of AI Systems Nobody Talks About. undefined vs undeclared, and how typeof behaves Switching from file-based jobs to NATS/Kafka in Rust without changing code io_uring Adventures: Rust Servers That Love Syscalls Why Agentic AI is Killing the Traditional Database The POUR principles of web accessibility for developers and designers Quantum Neural Network 3D — A Deep Dive into Interactive WebGL Visualization How To Install Caveman In Codex On macOS And Windows Automation Pipeline Reliability: Why Your Workflow Breaks When Nobody Is Watching I Built an 'Open World' AI Coding Agent — It Works From ANY Folder From Freelancing to Product: A Tech Service Company's SaaS Transformation China's AI Giants: Adding Tencent Hunyuan & ByteDance Doubao to AI University (74 Providers) On the Vibe Coders and Their Lies clerk: Auto-Summarize Your Claude Code Sessions AI Weekly — 2026/04/10–04/17 | The Model Lockdown Is Here, but the Toolchain Is the Real Battleground AI 週報 — 2026/04/10–2026/04/17 模型封鎖潮來了,但工具鏈才是真戰場 Maybe this is how Open-Source apps are born... 🚀 Fine-Tune LLMs with LoRA and QLoRA: 2026 Guide tRPC v11 + Next.js App Router: End-to-End Type Safety Without the Boilerplate ShadCN UI in 2026: Why I Stopped Installing Component Libraries and Started Owning My Components SaaS Billing in React Server Components: Stripe + Supabase Without a Single `useEffect` Join our DEV Weekend Challenge — $1,000 in Prizes Across TEN winners! Submissions Due April 20 at 6:59 AM UTC. Implementing FSRS Spaced Repetition in Flutter + Supabase — Adding Memory Science to an AI Learning App "I Texted My Localhost From the Train — Claude Code Fixed the Bug Before I Got Home" I Built a Sales Prep AI and It Went Deeper Than Expected Design to Code #2: One JSON, Eleven Outputs Solving the 100M-Row Problem: A Summary Table Pattern for High-Volume Push Notification Logs Flutter Web With Wasm: What Actually Changes For Developers I Built 50 Royalty-Free Soundtracks for My Side Project in a Weekend Using AI Music Generation The Vibe Coding Security Checklist: 7 Things to Check Before You Ship Stop Letting Googlebot Guess Fix Your React App's SEO Right Desconstruindo o Streaming do LinkedIn: Como Criar um Engine de Extração de Vídeo de Alta Performance com HLS e FFmpeg (EDA Part-1) EDA (Exploratory Data Analysis) Explained With Real Life — Why Looking at Your Data Is the Most Important Step in Machine Learning Brand Relationship Management at Scale: Our 4-Touch Outreach System for 200+ Brands Why String.fromEnvironment() Might Return an Empty String in Dart JGuardrails 1.0.0 — Hardening Java LLM Apps Against Jailbreaks, Toxicity, and Prompt Injection Plan and Schedule a Full Week of Threads Content From One Claude Conversation Coding Cat Oran Ep3, Five Tables Changed Everything Updated: BFF Pattern I'm done watching freelancers get buried by 200 proposals. So I'm building the alternative. This is my first post BFS Algorithm in Java Step by Step Tutorial with Examples Tracking LLM Pricing Monthly: An Open Dataset for 22 AI Models How We Measure Content ROI on a Comparison Site: Revenue Attribution Without Perfect Data Introducing Nova AI Ops: The AI-Native Operating System for SRE Teams I built a free desktop video downloader for Windows — Grabbit How Talkie OCR Helps Vision-Impaired & Dyslexic Users Read the World Around Them VRCFaceTracking安装和iPhone面捕配置教程,有bug Even CrowdStrike Can't See Your Agents The Automation Gold Rush: What n8n Workflows and Claude Are Opening Up for Developers Right Now
Running Async Python Inside Celery Is Harder Than You Think.
Kolade Fajimi · 2026-06-23 · via DEV Community

The problem is straightforward to state and surprisingly hard to solve correctly.

Celery workers are synchronous. Celery spawns prefork worker processes, and when a task arrives, it calls your task function like this: task_function(*args, **kwargs). It expects a return value. It blocks the worker thread until it gets one. It does not know or care that you wrote async def.

But modern Python services are async. FastAPI is async. SQLAlchemy 2.0 is async. httpx, aiohttp, asyncpg the entire interesting half of the ecosystem has gone async-first. The idea of maintaining two parallel code paths, one async for your web layer, one sync for your task layer is exactly the kind of thing that creates maintenance debt, copy-paste bugs, and the kind of divergence you only notice when something breaks in production.

So you want to write async def task functions and have them work inside a Celery worker. How hard can it be?

Harder than it looks.

Why asyncio.run() doesn't work

The first thing most people try:

def task_wrapper(*args, **kwargs):
    return asyncio.run(your_async_function(*args, **kwargs))

This works in isolation. It fails in production for a specific reason: asyncio.run() creates a new event loop, runs the coroutine to completion, then closes the loop. If there is already a running event loop on the current thread, and there frequently is, in test environments, in newer Celery versions, in signal handlers, it raises:

RuntimeError: This event loop is already running.

The fix most people find next is nest_asyncio:

import nest_asyncio
nest_asyncio.apply()
# now asyncio.run() "works" from inside a running loop

nest_asyncio patches the event loop to allow re-entrant calls. It works in simple cases. The subtle failure mode: re-entrant event loops change the execution order of scheduled callbacks and coroutines. Code that was safe under normal scheduling assumptions becomes non-deterministic under concurrent load. Bugs appear only at production concurrency, only under specific timing, and are nearly impossible to reproduce in development.

The prefork complication

Even if you solve the asyncio.run() problem, Celery's prefork concurrency model introduces a second failure that takes longer to diagnose because it manifests as infinite silence rather than an immediate error.

When Celery starts, it forks N worker processes from a single parent. After fork(), the child process inherits the parent's memory including any event loop objects that existed before the fork.

The problem: fork() does not copy threads. A Python asyncio.AbstractEventLoop is driven by a thread calling loop.run_forever(). After fork(), the child has the loop object but not the thread running it. The loop's internal state may indicate it was running; nothing is actually driving it. Any coroutine scheduled onto this loop hangs indefinitely.

@worker_process_init.connect
def bad_init(**kwargs):
    loop = asyncio.get_event_loop()
    # This loop was inherited from the parent.
    # The thread driving it died when the parent forked.
    # loop.is_running() → False.
    # Scheduling coroutines onto it produces no results and no errors.
    future = asyncio.run_coroutine_threadsafe(some_coro(), loop)
    future.result()  # blocks forever

This is the kind of bug that produces a zero-width failure window. The loop object exists and looks valid. No exception is raised. Work just never completes. I spent the better part of a day convinced the issue was in the Redis client before realizing the loop scheduled to drive it had died at fork time.

The solution: a persistent bridge loop per worker process

The correct approach is to create a brand-new event loop inside each forked worker process and start a dedicated daemon thread to drive it. The bridge loop is the only asyncio runtime in the worker process. All async work runs on it. Celery's synchronous worker threads never touch an event loop directly.

worker_loop: asyncio.AbstractEventLoop | None = None

@worker_process_init.connect
def init_worker_process(**kwargs):
    global worker_loop

    # Always create a fresh loop in the forked child.
    # Never reuse the inherited parent loop object.
    worker_loop = asyncio.new_event_loop()

    # A daemon thread drives the loop independently of Celery's
    # synchronous execution threads.
    t = threading.Thread(
        target=_run_event_loop,
        args=(worker_loop,),
        daemon=True
    )
    t.start()

def _run_event_loop(loop):
    asyncio.set_event_loop(loop)
    loop.run_forever()

Now the bridge is asyncio.run_coroutine_threadsafe. When Celery calls the synchronous task wrapper, the wrapper schedules the async orchestration coroutine onto the background loop and blocks the worker thread waiting for the result:

def wrapper(self, *args, **kwargs):
    async def _orchestrate():
        # Schema migration, idempotency check, Phoenix heartbeat,
        # OTel span setup, task execution, fence validation, DLQ quarantine.
        result = await your_async_task_function(*args, **kwargs)
        return result

    # Schedule the coroutine from this synchronous thread onto
    # the event loop running on the background thread.
    future = asyncio.run_coroutine_threadsafe(_orchestrate(), worker_loop)

    # Block the Celery worker thread here. All actual work happens
    # on the bridge loop thread.
    return future.result(timeout=300)

run_coroutine_threadsafe is the correct API for this pattern. It is thread-safe, it returns a concurrent.futures.Future (not an asyncio Future), and future.result() blocks without touching the event loop. The background loop thread does all the async I/O. The Celery worker thread just waits.

This solves both problems cleanly:

  • No asyncio.run() from inside a running loop. The loop lives on a different thread.
  • No inherited-but-dead loop. Each worker creates its own after fork.

The push/apush split

Dispatching tasks has its own version of this problem. Celery's send_task is synchronous and blocking, it opens a broker connection and writes a message. If you call it from inside an async FastAPI route handler, you block the event loop during a network round-trip.

This is why Relier has two dispatch methods:

# From async code (FastAPI, async Django):
receipt = await send_invoice.apush(invoice_id)

# From sync code (Flask routes, sync Django views, scripts):
receipt = send_invoice.push(invoice_id)

apush runs the blocking broker send in an executor so the async caller is never blocked:

async def apush(self, *args, **kwargs):
    # Admission check, schema wrapping, OTel context injection...

    loop = asyncio.get_running_loop()
    return await loop.run_in_executor(
        None,
        lambda: celery_app.send_task(
            self.name,
            args=(envelope,),
            queue=queue,
            task_id=task_id,
        ),
    )

push explicitly guards against being called from inside a running loop:

def push(self, *args, **kwargs):
    try:
        asyncio.get_running_loop()
    except RuntimeError:
        pass  # No running loop on this thread. Safe.
    else:
        raise RuntimeError(
            f"{self.name}.push() was called from inside a running event loop, "
            "where it would block and deadlock that loop. "
            f"Use `await {self.name}.apush(...)` instead."
        )

    # Inside a Celery worker: reuse the bridge loop.
    if worker_loop and worker_loop.is_running():
        future = asyncio.run_coroutine_threadsafe(
            self.apush(*args, **kwargs), worker_loop
        )
        return future.result(timeout=5.0)

    # Outside Celery (Flask route, script):
    return asyncio.run(self.apush(*args, **kwargs))

The error message in the RuntimeError matters. When someone calls push() from a FastAPI route handler, they get an actionable message telling them exactly what to do instead. Not a silent deadlock. Not a timeout with no context. A specific message at the exact moment the mistake is made.

The check itself, asyncio.get_running_loop() in a try/except RuntimeError is the canonical way to detect whether the current thread is running an event loop. It raises RuntimeError if no loop is running on this thread, which is the safe case for push().

Sync tasks in an async world

What about existing sync task functions? A codebase of def tasks shouldn't require a full rewrite to benefit from Relier's reliability stack.

Inside the orchestration coroutine, execution branches on whether the function is async:

if inspect.iscoroutinefunction(func):
    result = await func(*actual_args, **actual_kwargs)
else:
    result = await asyncio.to_thread(func, *actual_args, **actual_kwargs)

asyncio.to_thread runs the sync function in Python's default thread pool executor. The orchestration layer awaits it without blocking the bridge loop. All the async infrastructure, heartbeat refreshes, Phoenix registration, OTel span updates, fence validation keeps running concurrently on the bridge loop while the sync function runs on a thread pool thread.

The constraint is honest: two-tier timeouts (soft_timeout, hard_timeout) only work for async def tasks. A sync function running in asyncio.to_thread cannot be cooperatively cancelled from outside. Relier raises ValueError at decoration time if you pass timeout parameters to a sync task, rather than silently providing no protection:

@rl_task(soft_timeout=8, hard_timeout=10)  # ValueError at import time
def sync_task(data: str) -> dict:
    ...

# Fix: convert to async def, or remove the timeout parameters.

Failing loudly at decoration time is better than failing silently at runtime when the timeout fires and nothing happens.

Timeout enforcement without thread kills

Two-tier timeouts deserve their own explanation because they interact with the bridge loop in a non-obvious way.

When a task starts, Relier spawns two watcher coroutines as asyncio tasks alongside the actual work:

task_coro = asyncio.create_task(func(*args, **kwargs))

async def _soft_timeout_handler():
    await asyncio.sleep(float(soft))
    if not task_coro.done():
        # Fire the recovery hook. Task keeps running.
        if on_soft:
            await on_soft(ctx)

async def _hard_timeout_handler():
    await asyncio.sleep(float(hard))
    if not task_coro.done():
        task_coro.cancel()  # Delivers CancelledError at next await point.

soft_watcher = asyncio.create_task(_soft_timeout_handler())
hard_watcher = asyncio.create_task(_hard_timeout_handler())

done, pending = await asyncio.wait(
    {task_coro, hard_watcher},
    return_when=asyncio.FIRST_COMPLETED,
)

All three coroutines run concurrently on the bridge loop. The soft timeout fires and calls your recovery hook, where you can call ctx.set_partial(state) to checkpoint work in progress while the task keeps running. If the task doesn't finish before the hard deadline, task_coro.cancel() delivers asyncio.CancelledError at the task's next await point.

No thread kills. No SIGALRM. No OS-level signals. Pure cooperative asyncio cancellation. This matters for cleanup: CancelledError propagates through finally blocks. Resources get released. Partial state gets checkpointed. The task gets quarantined to the DLQ with its full payload and resurrection history. None of that happens with a hard OS kill.

The disposable loop case

One edge case worth knowing: outside a Celery worker in a CLI script, a management command, a test, there's no bridge loop. The task wrapper's loop resolution falls through to creating a fresh event loop just for that call:

def _get_worker_loop():
    # 1. Check for persistent worker bridge.
    if relier.tasks.app.worker_loop is not None:
        return relier.tasks.app.worker_loop

    # 2. Check for a running loop on this thread (test contexts).
    try:
        return asyncio.get_running_loop()
    except RuntimeError:
        pass

    # 3. Create a disposable loop for this one call.
    loop = asyncio.new_event_loop()
    asyncio.set_event_loop(loop)
    return loop

Disposable loops are cleaned up after the call: Redis connections are closed, the loop is stopped and closed, and asyncio.set_event_loop(None) clears the thread-local reference. The persistent worker_loop is specifically excluded from this cleanup path closing the bridge loop mid-execution would kill all in-flight tasks.

What I learned

The prefork problem is the kind of failure that shows up as "nothing happens" rather than an exception. You schedule coroutines, they don't run, no error surfaces. It took a day of debugging the wrong thing before I isolated it to the inherited-but-dead loop. The fix (create a fresh loop in worker_process_init) is obvious in retrospect. Getting there required understanding exactly what fork() does to threads.

asyncio.run_coroutine_threadsafe is underused. Most Python developers never need to cross a thread boundary into a running event loop, so the API is obscure. But for anything that marries a sync framework (Celery, Django ORM, WSGI in general) with async internals, it is the correct and safe way to do it. It appears in the Python docs in a single paragraph. It deserves more.

The two-method dispatch split (push/apush) is the right API surface even though it introduces surface area. The alternative, a single method that auto-detects the context and does the right thing sounds better but produces confusing failures when the auto-detection is wrong. The explicit split makes the contract clear. Async code always uses apush. Sync code always uses push. The guard in push() exists so that misuse produces a useful error immediately rather than a deadlock ten seconds later.

Cooperative timeout cancellation is better than OS-level signals for tasks that care about cleanup. The finally block guarantee is the part that matters: partial state can be persisted, connections can be closed, the DLQ entry gets written with everything needed to re-inspect or re-dispatch. An OS kill gives you none of that.

The whole bridge, bridge loop thread, run_coroutine_threadsafe, push/apush split, disposable loop cleanup is about 200 lines in app.py and decorator.py combined. The complexity is real but contained. Once the pattern is in place, every async def task function just works, without the task author knowing anything about the event loop infrastructure underneath.


GitHub: github.com/getrelier/relier

Docs: getrelier.github.io/relier

pip install relier