Menghasilkan dan Menggunakan Peristiwa Kafka Secara Tak Segerak
Integrasikan aiokafka dengan FastAPI untuk menerbitkan dan menggunakan peristiwa tanpa menyekat gelung peristiwa.
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
senddangetoneialah rutin tak segerak yang andaawait. - 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.statesupaya 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.
RecordMetadatayang 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_serializermenerima objek anda dan mesti mengembalikanbytes.- 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 futPerangkap 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.valuedengan cara yang sama seperti pengekodan di sisi pengeluar. - Anda juga boleh menghantar
value_deserializerkepada pembina pengguna. msg.timestampmengandungi 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 offsetKeserentakan 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_idyang 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 messageSemakan 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_waitmengesahkan penghantaran;sendmemaksimumkan 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.
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
- Menghasilkan dan Menggunakan Peristiwa Kafka Secara Tak Segerak
- Pendaftaran Skema dan Evolusi Kontrak Avro
- Corak Peti Keluar Transaksi
- Pengguna Idempoten dan Semantik Tepat Sekali