طراحی خط لوله ایونتمحور برای حذف خطاهای ثبت انبارداری زنجیره تامین توسن
پیادهسازی آپاچی کافکا با تضمین تحویل حداقل یکبار (At-least-once) و سیستم Idempotency جهت هماهنگسازی بلادرنگ ۱۲۰۰ شعبه توزیع کالا.
سهراب حسینیمهندس داده و اتوماسیون
زنجیره تامین توسن با ۱۲۰۰ شعبه توزیع کالا، روزانه بیش از ۲ میلیون رویداد ثبت انبار میفرستد. سامانه قبلی — یک REST API متمرکز روی یک پایگاهداده — در ساعات اوج با خطاهای تداخل (Race Condition) کسری و اضافهمجوزی موجودی دستوپنجه نرم میکرد. بازطراحی ایونتمحور این مشکل را ریشهای حل کرد.
مدل رویداد بهجای فراخوانی مستقیم
بهجای اینکه شعبه مستقیم موجودی را «کم» کند، یک رویداد StockAdjustmentRequested ثبت میکند و همه چیز از آنجا جریان مییابد. سه اصل برقرار است:
- هر تغییر، یک رویداد با شناسه idempotent: ارسال تکراری همان تراکنش، اثر تکراری ندارد.
- ترتیب per-aggregate: رویدادهای هر SKU در یک پارتیشن Kafka با همان کلید مرتب میمانند.
- جدا بودن نوشتن از خواندن: نمودارهای موجودی از روی رویدادها ساخته میشوند، نه از صف عملیات.
وقتی تراکنش را به رویداد تبدیل میکنید، دیگر با قطعی شبکه نمیجنگید؛ فقط با زمان — و زمان را میتوان با بازپخش مدیریت کرد.
تضمین At-least-once و نقش Idempotency
Kafka تحویل حداقل-یکبار میدهد؛ یعنی مصرفکننده باید هر پیام را طوری بنویسد که تکرارش بیاثر باشد. الگوی ما یک جدول تشخیص تکرار با کلید یکتای رویداد است:
-- مصرفکننده: ثبت فقط-یکبار هر رویداد
INSERT INTO processed_event (event_id, consumer, processed_at)
VALUES ($1, $2, now())
ON CONFLICT (event_id, consumer) DO NOTHING;
-- فقط اگر ردیف واقعاً درج شد، اثر رویداد اعمال میشود
INSERT INTO stock_movement (sku, branch_id, delta, event_id)
SELECT $3, $4, $5, $1
WHERE NOT EXISTS (
SELECT 1 FROM stock_movement WHERE event_id = $1
);
مقابله با پیامهای خارج از ترتیب
گاهی مصرفکننده رویدادِ جدیدتر را قبل از قدیمیتر میبیند (مثلاً بعد از seek). راهحل: تاریخ اعتبار رویداد را نگه میداریم و اگر قدیمیتر از آخرین وضعیت ثبتشده باشد، فقط در لاگ رد میشود.
def apply(event):
current = load_state(event.sku, event.branch)
if current is not None and event.occurred_at <= current.updated_at:
log.discard(event.id, reason="out-of-order")
return
write_state(event)
نتایج پس از استقرار
| معیار | قبل (REST متمرکز) | بعد (ایونتمحور) |
|---|---|---|
| خطای تداخل موجودی | روزانه ~۴۰۰ مورد | صفر |
| تأخیر هماهنگسازی شعب | ۵ تا ۱۵ دقیقه | زیر ۳ ثانیه |
| افت سرویس در اوج | ۲ رویداد در ماه | بدون رویداد |

نتیجهگیری
هیچیک از این تغییرات «فناوری برای فناوری» نبود: هر تصمیم مستقیماً یک خطای ثبتشده در تیکتهای پشتیبانی شعب را هدف میگرفت. اگر سامانه شما هم با تداخل نوشتنهای همزمان دستوپنجه نرم میکند، اول مدل رویداد را بکشید؛ تازه بعد درباره کافکا فکر کنید.
برچسبها
سهراب حسینی
مهندس داده و اتوماسیوندرباره نویسنده
طراحی پایپلاینهای ایونتمحور، صفهای توزیعشده و سامانههای پایش برای سازمانهای با دادههای تراکنشی سنگین.
مقالات نویسنده

