Kem Intensif Pembangunan Bahagian Belakang FastAPI · Pelajaran

Menghasilkan dan Menggunakan Peristiwa Kafka Secara Tak Segerak

Integrasikan aiokafka dengan FastAPI untuk menerbitkan dan menggunakan peristiwa tanpa menyekat gelung peristiwa.

Pelajaran 1 daripada 413 langkah

Menghasilkan dan Menggunakan Peristiwa Kafka Secara Tak Segerak ialah pelajaran Kem Intensif Pembangunan Bahagian Belakang FastAPI percuma di CoddyKit. Ini ialah pelajaran 1 daripada 4. Sebanyak 3 pelajaran dalam laluan pembelajaran ini boleh dibaca sepenuhnya secara percuma — selepas itu, CoddyKit PRO membuka akses kepada semua pelajaran, serta latihan praktikal dengan penyunting kod terbina dalam dan tutor kecerdasan buatan yang tersedia 24/7. Pelajaran ini merupakan sebahagian daripada laluan pembelajaran Kem Intensif Pembangunan Bahagian Belakang FastAPI, dan kemajuan anda disegerakkan merentas web serta aplikasi CoddyKit. Kursus Kem Intensif Pembangunan Bahagian Belakang FastAPI merangkumi sejumlah 4 pelajaran.

Mengapa Kafka Tak Segerak dalam FastAPI

FastAPI berjalan pada gelung acara asyncio. Jika anda menerbitkan atau meninjau Kafka dengan klien menyekat (seperti pustaka kafka-python standard), setiap panggilan rangkaian membekukan keseluruhan gelung dan menghentikan semua permintaan serentak.

  • aiokafka ialah klien Kafka asyncio asli yang tidak pernah menyekat gelung.
  • Operasi send dan getone ialah rutin tak segerak yang anda await.
  • Hal ini membolehkan satu pekerja mengendalikan ribuan permintaan yang sedang berjalan sementara I/O Kafka menunggu.

Dalam pelajaran ini, anda akan menyambungkan AIOKafkaProducer dan AIOKafkaConsumer kepada aplikasi FastAPI dengan cara yang betul.

Kitar Hayat Pengeluar dengan lifespan

Pengeluar mengekalkan sambungan TCP dan tugas penghantar latar belakang. Anda mesti menjalankan start() sekali sahaja ketika aplikasi dimulakan dan menjalankan stop() ketika penutupan — jangan sekali-kali bagi setiap permintaan. Cara moden FastAPI ialah pengurus konteks lifespan.

  • await producer.start() membuka sambungan dan gelung penghantar.
  • await producer.stop() menghantar semua kelompok yang menunggu dan menutup sambungan dengan kemas.
  • Simpan pengeluar dalam app.state supaya laluan boleh mencapainya.
from contextlib import asynccontextmanager
from fastapi import FastAPI
from aiokafka import AIOKafkaProducer


@asynccontextmanager
async def lifespan(app: FastAPI):
    producer = AIOKafkaProducer(
        bootstrap_servers="localhost:9092",
        enable_idempotence=True,
    )
    await producer.start()
    app.state.producer = producer
    try:
        yield
    finally:
        await producer.stop()


app = FastAPI(lifespan=lifespan)

Menerbitkan Peristiwa daripada Laluan

Dalam laluan, ambil pengeluar yang dikongsi dan gunakan await producer.send_and_wait(...). Panggilan send_and_wait kembali setelah broker mengakui rekod tersebut, lalu memberikan kawalan tekanan balik dan pengesahan penghantaran.

  • Kunci dan nilai Kafka ialah bait — kodkan JSON sendiri atau hantarkan pensiri.
  • RecordMetadata yang dikembalikan memberitahu anda tentang petak dan ofset.
  • Gunakan kunci mesej untuk menjamin susunan bagi setiap entiti merentasi petak.
import json
from fastapi import FastAPI, Request

app = FastAPI()


@app.post("/orders")
async def create_order(payload: dict, request: Request):
    producer = request.app.state.producer
    value = json.dumps(payload).encode("utf-8")
    key = str(payload["order_id"]).encode("utf-8")
    meta = await producer.send_and_wait(
        "orders", value=value, key=key
    )
    return {"partition": meta.partition, "offset": meta.offset}

Pensiri berbanding Pengekodan Manual

Daripada memanggil json.dumps(...).encode() pada setiap penghantaran, anda boleh memberikan aiokafka value_serializer dan key_serializer. Pengeluar menggunakannya secara automatik, jadi laluan boleh menghantar objek Python biasa.

  • value_serializer menerima objek anda dan mesti mengembalikan bytes.
  • Hal ini memusatkan pengekodan dan mengelakkan pengulangan merentasi laluan.
  • Dalam pasukan pengeluaran, JSON sering digantikan dengan Avro atau Protobuf melalui pendaftaran skema.
import json
from aiokafka import AIOKafkaProducer

producer = AIOKafkaProducer(
    bootstrap_servers="localhost:9092",
    value_serializer=lambda v: json.dumps(v).encode("utf-8"),
    key_serializer=lambda k: str(k).encode("utf-8"),
)

# Now routes can send native objects:
# await producer.send_and_wait("orders", value={"id": 7}, key=7)

send berbanding send_and_wait

Dua kaedah pengeluar, dua pertukaran:

  • send() mengembalikan future serta-merta dan membolehkan aiokafka mengumpulkan rekod di latar belakang — daya pemprosesan tinggi, tetapi anda masih belum tahu sama ada penghantaran berjaya.
  • send_and_wait() menunggu pengakuan broker — lebih perlahan bagi setiap panggilan, tetapi memberikan hasil dan mendedahkan ralat dengan segera.

Bagi pengendali permintaan yang memerlukan pengesahan daripada klien, utamakan send_and_wait. Bagi penghantaran pukal secara terus tanpa menunggu, gunakan send dan secara pilihan tunggu future tersebut kemudian.

# Fire many records fast, then wait once for all of them
futures = []
for item in batch:
    fut = await producer.send("events", value=item)
    futures.append(fut)

# Awaiting the futures surfaces any delivery errors
for fut in futures:
    record_meta = await fut

Perangkap Menyekat yang Perlu Dielakkan

Kesilapan yang paling lazim ialah mencampurkan klien segerak ke dalam aplikasi tak segerak. Di bawah ialah demonstrasi kendiri tentang sebab menyekat gelung mendatangkan masalah: time.sleep yang menyekat di dalam rutin tak segerak menghentikan segala-galanya, manakala asyncio.sleep menyerahkan kawalan.

Jalankan kod ini dan perhatikan bahawa versi yang menggunakan penantian membolehkan kedua-dua tugas bertindih — tepat seperti tingkah laku yang diberikan oleh aiokafka untuk I/O rangkaian sebenar.

import asyncio
import time


async def good_io(name):
    await asyncio.sleep(0.2)  # yields the loop
    print(f"{name} done at {time.strftime('%X')}")


async def main():
    start = time.perf_counter()
    await asyncio.gather(good_io("A"), good_io("B"))
    print(f"both finished in {time.perf_counter() - start:.2f}s")


asyncio.run(main())

Pengguna sebagai Tugas Latar Belakang

Aplikasi FastAPI menyediakan HTTP, tetapi pengguna Kafka mesti meninjau secara berterusan. Corak yang kemas ialah memulakan pengguna dalam lifespan dan menjalankan gelung tinjauannya sebagai tugas latar belakang asyncio, lalu membatalkannya ketika penutupan.

  • asyncio.create_task(...) melancarkan gelung tanpa menyekat permulaan.
  • Semasa penutupan, batalkan tugas itu, kemudian gunakan await consumer.stop().
  • Sentiasa bungkus badan gelung supaya satu mesej yang bermasalah tidak mematikan pengguna.
import asyncio
from contextlib import asynccontextmanager
from fastapi import FastAPI
from aiokafka import AIOKafkaConsumer


async def consume(consumer: AIOKafkaConsumer):
    async for msg in consumer:
        try:
            handle(msg.value)
        except Exception as exc:
            print("handler failed:", exc)


@asynccontextmanager
async def lifespan(app: FastAPI):
    consumer = AIOKafkaConsumer(
        "orders",
        bootstrap_servers="localhost:9092",
        group_id="order-workers",
    )
    await consumer.start()
    task = asyncio.create_task(consume(consumer))
    try:
        yield
    finally:
        task.cancel()
        await consumer.stop()


app = FastAPI(lifespan=lifespan)

Mengulangi Mesej dan Menyahkod

AIOKafkaConsumer ialah lelaran tak segerak: async for msg in consumer menunggu setiap rekod baharu. Setiap msg mendedahkan topic, partition, offset, key, dan value sebagai bait.

  • Nyahkod msg.value dengan cara yang sama seperti pengekodan di sisi pengeluar.
  • Anda juga boleh menghantar value_deserializer kepada pembina pengguna.
  • msg.timestamp mengandungi cap masa broker atau pengeluar untuk metrik kependaman.
import json
from aiokafka import AIOKafkaConsumer

consumer = AIOKafkaConsumer(
    "orders",
    bootstrap_servers="localhost:9092",
    group_id="order-workers",
    value_deserializer=lambda b: json.loads(b.decode("utf-8")),
)


async def run():
    async for msg in consumer:
        event = msg.value  # already a dict
        print(event["order_id"], "at offset", msg.offset)

Komit Ofset: Penghantaran Sekurang-kurangnya Sekali

Komit automatik (lalai) melakukan komit ofset berdasarkan pemasa, yang boleh kehilangan mesej jika pekerja ranap selepas komit tetapi sebelum pemprosesan. Untuk pemprosesan yang boleh dipercayai, tetapkan enable_auto_commit=False dan lakukan komit selepas anda selesai mengendalikan mesej.

  • Komit manual memberikan semantik sekurang-kurangnya sekali — ranap sistem akan memainkan semula mesej terakhir yang belum dikomit.
  • Oleh sebab mesej boleh berulang, pengendali anda mestilah idempoten.
  • Lakukan komit dalam kelompok kecil untuk mengimbangi daya pemprosesan dengan kos main semula.
consumer = AIOKafkaConsumer(
    "orders",
    bootstrap_servers="localhost:9092",
    group_id="order-workers",
    enable_auto_commit=False,
    auto_offset_reset="earliest",
)


async def run():
    async for msg in consumer:
        await process(msg.value)   # do the work first
        await consumer.commit()    # then advance the offset

Keserentakan dan Susunan Petak

Dalam satu petak, Kafka mengekalkan susunan, dan gelung async-for memproses rekod tersebut secara berjujukan. Untuk membuat penskalaan, anda mempunyai dua pilihan:

  • Lebih banyak petak + lebih banyak pengguna dalam group_id yang sama — Kafka menetapkan petak merentasi tika secara automatik.
  • Keserentakan terhad bagi setiap pekerja dengan semaphore, apabila kerja bagi setiap mesej banyak menggunakan I/O dan susunan ketat tidak diperlukan.

Jika susunan mengikut entiti penting, kekalkan entiti itu dalam satu petak melalui kunci mesej yang stabil dan proses petaknya secara bersiri.

import asyncio

sem = asyncio.Semaphore(10)


async def handle_bounded(value):
    async with sem:
        await do_async_work(value)


async def run(consumer):
    async for msg in consumer:
        # Schedule work without blocking the poll loop
        asyncio.create_task(handle_bounded(msg.value))

Penutupan Lancar dan Pengendalian Ralat

Penerapan yang kukuh melakukan pembersihan supaya data yang sedang diproses tidak hilang:

  • Semasa penutupan, batalkan tugas pengguna dan gunakan await consumer.stop() — tindakan ini melakukan komit ofset (jika komit automatik digunakan) dan meninggalkan kumpulan dengan kemas supaya pengimbangan semula berlaku dengan pantas.
  • await producer.stop() menghantar rekod yang disimpan dalam penimbal sebelum menutup.
  • Bungkus pengendalian setiap mesej dalam try/except dan halakan mesej rosak ke topik surat mati dan bukannya mematikan gelung.

Jangan sekali-kali panggil start()/stop() dalam pengendali permintaan — tindakan itu mengganggu sambungan dan merosakkan keahlian kumpulan pengguna.

async def consume(consumer, dlq_producer):
    async for msg in consumer:
        try:
            await process(msg.value)
            await consumer.commit()
        except PermanentError:
            await dlq_producer.send_and_wait(
                "orders.dlq", value=msg.value, key=msg.key
            )
            await consumer.commit()  # skip the poison message

Semakan Pantas

Anda memerlukan pemprosesan yang boleh dipercayai, dengan syarat ranap pekerja tidak boleh menggugurkan peristiwa pesanan secara senyap. Konfigurasi pengguna yang manakah paling menyokong keperluan ini?

Imbas Kembali

Anda telah mengintegrasikan Kafka ke dalam FastAPI tanpa menyekat gelung peristiwa:

  • aiokafka menyediakan klien pengeluar dan pengguna yang boleh ditunggu, asli untuk asyncio.
  • Urus pengeluar dan pengguna dalam konteks jangka hayat: start() semasa but, stop() semasa penutupan — jangan sekali-kali bagi setiap permintaan.
  • send_and_wait mengesahkan penghantaran; send memaksimumkan daya pemprosesan.
  • Jalankan gelung tinjauan pengguna sebagai tugas latar belakang asyncio dan lakukan lelaran dengan async for.
  • Lumpuhkan komit automatik dan lakukan komit selepas pemprosesan untuk penghantaran sekurang-kurangnya sekali, sambil memastikan pengendali idempoten.
  • Skalakan dengan lebih banyak pembahagian dan pengguna dalam satu kumpulan, kekalkan susunan bagi setiap entiti melalui kunci mesej, dan lakukan penutupan dengan lancar menggunakan topik surat mati untuk mesej bermasalah.
Percuma untuk bermula

Pelajari Kem Intensif Pembangunan Bahagian Belakang FastAPI dengan tutor kecerdasan buatan — percuma

Tulis dan jalankan kod sebenar dalam pelayar anda, dapatkan bantuan segera daripada tutor kecerdasan buatan yang tersedia 24/7, dan sambung semula dari tempat anda berhenti di web atau dalam aplikasi.

Kursus
21
Pelajaran
84

Soalan Lazim

Adakah pelajaran “Menghasilkan dan Menggunakan Peristiwa Kafka Secara Tak Segerak” percuma?

Ya — sebanyak 3 pelajaran dalam laluan pembelajaran Kem Intensif Pembangunan Bahagian Belakang FastAPI, termasuk “Menghasilkan dan Menggunakan Peristiwa Kafka Secara Tak Segerak”, boleh dibaca sepenuhnya secara percuma di web ini. Selepas itu, CoddyKit PRO membuka akses kepada semua pelajaran, serta latihan interaktif dengan penyunting kod terbina dalam dan tutor kecerdasan buatan yang tersedia 24/7. Kursus Kem Intensif Pembangunan Bahagian Belakang FastAPI merangkumi sejumlah 4 pelajaran.

Apakah yang akan saya pelajari dalam “Menghasilkan dan Menggunakan Peristiwa Kafka Secara Tak Segerak”?

Integrasikan aiokafka dengan FastAPI untuk menerbitkan dan menggunakan peristiwa tanpa menyekat gelung peristiwa. Anda berlatih Kem Intensif Pembangunan Bahagian Belakang FastAPI menggunakan kod praktikal yang dijalankan terus dalam pelayar, manakala tutor kecerdasan buatan 24/7 menjawab soalan anda semasa anda mengikuti pelajaran.

Adakah saya memerlukan pengalaman untuk memulakan Kem Intensif Pembangunan Bahagian Belakang FastAPI?

Tiada pengalaman terdahulu diperlukan. Pembelajaran Kem Intensif Pembangunan Bahagian Belakang FastAPI di CoddyKit disusun untuk pelajar daripada peringkat pemula hingga lanjutan, jadi anda boleh bermula di sini atau dari awal dan belajar mengikut kadar anda sendiri. Ini ialah pelajaran 1 daripada 4.

Berapa lamakah pelajaran “Menghasilkan dan Menggunakan Peristiwa Kafka Secara Tak Segerak” diambil?

Kebanyakan pelajaran CoddyKit mengambil masa kira-kira 5–10 minit. Setiap pelajaran ringkas dan interaktif, jadi anda boleh membuat kemajuan secara berterusan dan menyambung tepat dari tempat anda berhenti di web atau aplikasi.

Bolehkah saya menulis dan menjalankan kod dalam pelajaran Kem Intensif Pembangunan Bahagian Belakang FastAPI ini?

Ya. Setiap pelajaran Kem Intensif Pembangunan Bahagian Belakang FastAPI menyertakan penyunting kod terbina dalam, jadi anda boleh menulis dan menjalankan kod sebenar terus dalam pelayar serta menerima maklum balas kecerdasan buatan serta-merta — tanpa memerlukan persediaan setempat.

Semua pelajaran dalam kursus ini

  1. Menghasilkan dan Menggunakan Peristiwa Kafka Secara Tak Segerak
  2. Pendaftaran Skema dan Evolusi Kontrak Avro
  3. Corak Peti Keluar Transaksi
  4. Pengguna Idempoten dan Semantik Tepat Sekali
← Kembali ke Kem Intensif Pembangunan Bahagian Belakang FastAPI