Pool worker asynchronous sering dimulai dengan pola sederhana: producer memasukkan job ke asyncio.Queue, consumer melakukan loop pada get(), dan aplikasi menunggu join() sebelum keluar.

Bagian yang canggung adalah shutdown. Desain lama biasanya memasukkan satu nilai sentinel ke queue untuk setiap worker, membatalkan consumer setelah join(), atau mempertahankan stop event terpisah. Setiap pendekatan dapat bekerja, tetapi masing-masing menambahkan protokol kedua di samping queue itu sendiri.

Python 3.13 menambahkan asyncio.Queue.shutdown() dan exception asyncio.QueueShutDown. Keduanya memungkinkan queue merepresentasikan lifecycle-nya sendiri: terbuka untuk producer, dalam proses shutdown sementara pekerjaan yang ada dikuras, lalu akhirnya tertutup bagi consumer.

Hal ini membuat graceful shutdown lebih mudah dipahami, tetapi hanya jika task_done() dan join() tetap mempertahankan makna yang dimaksudkan.

Mulai dari lifecycle queue

asyncio.Queue normal menerima pemanggilan put() dan memungkinkan consumer menunggu di get(). Bounded queue juga menyediakan backpressure: ketika mencapai maxsize, await queue.put(item) menunggu kapasitas tersedia.

Memanggil:

queue.shutdown()

mengubah kontrak tersebut. Queue tidak lagi dapat bertambah. Pemanggilan put() berikutnya melempar QueueShutDown, dan producer yang sudah terblokir di put() dibangunkan lalu menerima exception yang sama.

Dengan default immediate=False, item yang sudah ada di queue tetap tersedia. Consumer dapat mengurasnya secara normal. Setelah queue kosong, get() melempar QueueShutDown.

Hal itu memberi worker kondisi keluar yang natural:

import asyncio


async def worker(queue: asyncio.Queue[str]) -> None:
    while True:
        try:
            item = await queue.get()
        except asyncio.QueueShutDown:
            return

        try:
            await process(item)
        finally:
            queue.task_done()

Block finally penting. Setelah get() berhasil, item tersebut berkontribusi pada unfinished-task count milik queue sampai tepat satu task_done() mengakuinya.

Kuras pool worker secara graceful

Pool worker kecil dapat menggunakan shutdown tanpa objek sentinel:

import asyncio


async def process(item: str) -> None:
    await asyncio.sleep(0.05)
    print(item)


async def worker(name: str, queue: asyncio.Queue[str]) -> None:
    while True:
        try:
            item = await queue.get()
        except asyncio.QueueShutDown:
            return

        try:
            await process(item)
        finally:
            queue.task_done()


async def main() -> None:
    queue: asyncio.Queue[str] = asyncio.Queue(maxsize=100)

    async with asyncio.TaskGroup() as group:
        for number in range(4):
            group.create_task(worker(f"worker-{number}", queue))

        for number in range(1_000):
            await queue.put(f"job-{number}")

        queue.shutdown()
        await queue.join()


asyncio.run(main())

Urutannya disengaja:

  1. Hentikan penerimaan pekerjaan baru dengan shutdown().
  2. Biarkan worker memproses semua pekerjaan yang sudah diterima.
  3. Tunggu acknowledgement dengan join().
  4. Biarkan setiap worker menerima QueueShutDown setelah queue terkuras lalu return.

Tidak diperlukan protokol sentinel yang bergantung pada jumlah worker.

Perlakukan join() sebagai acknowledgement barrier

join() tidak menunggu queue menjadi kosong. Ia menunggu unfinished-task count mencapai nol.

Setiap put() yang berhasil menaikkan count tersebut. Setiap task_done() yang sesuai menurunkannya. Perbedaan ini penting karena consumer dapat mengeluarkan item dari queue jauh sebelum selesai memproses item tersebut.

Pertimbangkan database writer:

item = await queue.get()
try:
    await write_to_database(item)
finally:
    queue.task_done()

Memanggil task_done() segera setelah get() akan membuat join() dapat return ketika database write masih berjalan. Saat deployment atau process shutdown, hal itu dapat mengubah coordination primitive menjadi sinyal durability yang palsu.

Invariant yang berguna adalah:

successful get() -> processing finishes or fails -> exactly one task_done()

Jika kegagalan berarti item harus di-retry secara durable, desain retry policy tersebut secara terpisah. task_done() hanya memperhitungkan item pada in-memory queue; ia tidak membuktikan keberhasilan pada level bisnis.

Tangani kegagalan processing secara eksplisit

Pola finally menjaga accounting queue tetap benar, tetapi dapat menyembunyikan pertanyaan desain penting: apakah satu job yang gagal harus menghentikan worker pool?

Jika process() melempar exception dan worker membiarkannya keluar, TaskGroup akan membatalkan sibling task. Itu mungkin tepat untuk pelanggaran invariant, tetapi biasanya terlalu agresif untuk error per-item yang memang diperkirakan.

Tangani kegagalan yang dapat dipulihkan di dalam worker:

async def worker(queue, failures):
    while True:
        try:
            item = await queue.get()
        except asyncio.QueueShutDown:
            return

        try:
            await process(item)
        except ExpectedJobError as exc:
            await failures.put((item, exc))
        finally:
            queue.task_done()

Untuk job yang tidak boleh hilang, in-memory queue sebaiknya bukan satu-satunya source of truth. Persist job atau retry state-nya sebelum mengakui sistem durable apa pun yang memilikinya.

Pahami apa yang terjadi pada producer yang terblokir

Shutdown bukan hanya fitur consumer.

Misalkan bounded queue penuh dan beberapa producer sedang menunggu di sini:

await queue.put(item)

Setelah task lain memanggil queue.shutdown(), producer yang terblokir tersebut bangun dan melempar QueueShutDown. Ini berguna karena shutdown tidak perlu menunggu kapasitas hanya untuk memberi tahu producer bahwa pekerjaan baru tidak lagi diterima.

Producer harus menentukan apakah exception tersebut merupakan lifecycle control yang diharapkan atau sebuah error:

async def producer(queue, source):
    async for item in source:
        try:
            await queue.put(item)
        except asyncio.QueueShutDown:
            return

Jangan secara membabi buta menangkap Exception di seluruh producer lalu melanjutkan. Setelah shutdown dimulai, mencoba ulang put() pada queue yang sama tidak dapat membukanya kembali.

Pilih graceful shutdown daripada immediate shutdown

Queue.shutdown() juga menerima immediate=True:

queue.shutdown(immediate=True)

Ini adalah operasi yang berbeda, bukan sekadar graceful shutdown yang lebih cepat. Queue dikuras segera. Getter yang terblokir dibangunkan dan melempar QueueShutDown karena tidak ada lagi item dalam queue untuk diambil.

Yang paling penting, immediate shutdown dapat melepaskan join() meskipun pekerjaan yang ada di queue tidak pernah diproses. Hal itu dengan sengaja melanggar invariant normal join().

Gunakan immediate shutdown hanya ketika membuang queued work dapat diterima atau ketika kegagalan pada level yang lebih tinggi sudah membuat draining normal tidak mungkin dilakukan.

Sebagai contoh, aplikasi dapat memilihnya setelah fatal dependency failure:

try:
    await run_pipeline(queue)
except FatalPipelineError:
    queue.shutdown(immediate=True)
    raise

Jangan menginterpretasikan join() yang return setelah immediate shutdown sebagai bukti bahwa semua job yang diterima sudah selesai.

Pisahkan graceful shutdown dari cancellation

Queue shutdown dan task cancellation menyelesaikan masalah yang berbeda.

queue.shutdown() mengubah apa yang dapat dilakukan producer dan consumer terhadap queue. Cancellation menginterupsi coroutine pada await point.

Untuk terminasi service normal, queue shutdown sering menjadi langkah pertama yang lebih bersih karena worker menyelesaikan job yang sudah mereka miliki dan menguras job yang sudah diterima sebelum keluar.

Cancellation tetap berguna ketika deadline habis:

queue.shutdown()

try:
    async with asyncio.timeout(10):
        await queue.join()
except TimeoutError:
    queue.shutdown(immediate=True)
    raise

Struktur task di sekitarnya kemudian dapat membatalkan worker jika aplikasi tidak dapat menunggu lebih lama. Ingat bahwa cancellation dapat menginterupsi process(item) setelah get() berhasil. Worker sebaiknya menggunakan finally untuk accounting queue dan membuat external side effect idempotent ketika retry dimungkinkan.

Hindari mencampur protokol shutdown tanpa alasan jelas

Queue yang menggunakan shutdown() biasanya tidak juga membutuhkan None, object sentinel, dan stop event.

Beberapa protokol sekaligus menciptakan state yang ambigu. Producer dapat memasukkan sentinel sebelum producer lain selesai. Worker dapat keluar karena event sementara item masih tersisa dalam queue. Sentinel juga dapat memakan kapasitas bounded queue dan membutuhkan pengetahuan tentang berapa banyak consumer yang harus menerimanya.

Tetap ada kasus ketika marker pada level data bermakna. Sebagai contoh, stream dapat memuat explicit partition-end record yang merupakan bagian dari business protocol. Jaga marker tersebut tetap terpisah dari queue lifecycle control.

Koordinasikan beberapa producer sebelum menutup admission

Task yang memanggil shutdown() harus mengetahui bahwa tidak ada producer yang sah yang masih perlu mengirim pekerjaan baru.

Jika beberapa producer berjalan secara concurrent, jangan biarkan producer pertama yang selesai menutup shared queue. Sebagai gantinya, koordinasikan penyelesaian producer pada level yang lebih tinggi:

async def produce_all(queue, sources):
    async with asyncio.TaskGroup() as group:
        for source in sources:
            group.create_task(produce(queue, source))

    queue.shutdown()

Sekarang shutdown berarti semua producer sudah selesai, bukan hanya satu producer.

Untuk service yang berjalan lama, trigger-nya dapat berupa server lifecycle event. Prinsip yang sama berlaku: pertama hentikan source yang dapat membuat job baru, lalu tutup admission queue, kemudian kuras pekerjaan yang sudah diterima.

Pertahankan backpressure selama operasi normal

Dukungan shutdown tidak menggantikan capacity planning.

Unbounded queue dapat menyerap burst sementara, tetapi ketidakseimbangan producer-consumer yang berkelanjutan berubah menjadi pertumbuhan memory. Berikan queue maxsize yang terbatas ketika producer dapat menunggu dengan aman:

queue = asyncio.Queue(maxsize=500)

Ini membatasi queued item, bukan total penggunaan resource. Worker dapat memegang active job, dan setiap item dapat mereferensikan object besar. Pilih limit berdasarkan pengukuran workload, bukan memperlakukan item count sebagai memory limit.

Shutdown terintegrasi dengan baik bersama bounded queue karena producer yang terblokir dilepas secara eksplisit dengan QueueShutDown, alih-alih tetap macet di belakang queue yang tidak lagi ingin dikuras aplikasi untuk submission baru.

Jangan bagikan asyncio.Queue antar-thread

asyncio.Queue dirancang untuk async code dan tidak thread-safe. Jika thread harus bertukar pekerjaan, gunakan primitive yang thread-safe seperti queue.Queue, atau lewati event-loop boundary dengan mekanisme scheduling thread-safe yang sesuai.

Tipe queue dengan nama serupa memiliki konsep terkait, tetapi jangan perlakukan mereka sebagai synchronization object yang dapat saling dipertukarkan.

Uji transisi lifecycle, bukan hanya output happy path

Bug pada queue shutdown adalah timing bug, jadi test harus melatih state di sekitar transisi.

Kasus yang berguna meliputi:

  • shutdown ketika queue kosong;
  • shutdown ketika item sedang menunggu;
  • shutdown ketika worker memproses item terakhir;
  • shutdown ketika producer terblokir pada bounded queue yang penuh;
  • worker melempar exception sebelum task_done() biasanya dijalankan;
  • graceful shutdown diikuti join();
  • immediate shutdown dengan unfinished work;
  • percobaan put() berulang setelah shutdown.

Test terfokus dapat memverifikasi bahwa producer yang terblokir dilepas:

async def test_shutdown_releases_blocked_producer():
    queue = asyncio.Queue(maxsize=1)
    await queue.put("first")

    blocked = asyncio.create_task(queue.put("second"))
    await asyncio.sleep(0)

    queue.shutdown()

    try:
        await blocked
    except asyncio.QueueShutDown:
        pass
    else:
        raise AssertionError("producer should observe shutdown")

Uji juga jaminan pada level aplikasi. Jika join() yang berhasil seharusnya berarti semua database write sudah committed, injeksikan slow write dan failure untuk membuktikan bahwa task_done() terjadi pada boundary yang tepat.

Rencanakan kompatibilitas dengan sengaja

asyncio.Queue.shutdown() dan QueueShutDown ditambahkan pada Python 3.13. Kode yang harus berjalan pada Python 3.12 atau lebih lama tidak dapat menggunakannya secara langsung.

Untuk library yang mendukung runtime lama, pertahankan strategi sentinel atau cancellation yang sudah mapan sampai versi minimum Python yang didukung mencapai 3.13. Hindari compatibility shim parsial yang meniru nama method tanpa mereproduksi semantik wake-up dan unfinished-task; perilaku concurrency adalah bagian dari kontrak API.

Aplikasi yang sudah distandardisasi pada Python 3.13 atau lebih baru dapat menyederhanakan kode lifecycle worker dengan menjadikan queue itu sendiri sebagai shutdown boundary.

Bangun shutdown di sekitar invariant yang eksplisit

Desain yang paling andal bukan yang memiliki baris paling sedikit. Desain yang paling andal adalah yang state-nya memiliki makna jelas.

Selama operasi normal, bounded put() menyediakan backpressure. Selama graceful shutdown, put baru gagal sementara pekerjaan yang sudah diterima tetap dapat dikuras. Setiap item yang diambil diakui tepat satu kali setelah processing. join() berarti semua pekerjaan yang diterima sudah diakui. Worker keluar ketika queue yang kosong dan sudah shut down melempar QueueShutDown.

Simpan immediate shutdown untuk kasus ketika Anda memang sengaja meninggalkan jaminan completion tersebut.

Dengan invariant tersebut, asyncio.Queue.shutdown() menghapus banyak signaling buatan tangan yang sebelumnya dibutuhkan asynchronous worker pool, sekaligus membuat jalur shutdown lebih mudah diuji dan dijelaskan.