Data Engineering

تسلط بر خاصیت تکرارپذیری و معنای دقیقاً یک‌بار در آپاچی کافکا

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

چالش تکراری‌ها در سیستم‌های توزیع‌شده

در معماری‌های توزیع‌شده، خرابی‌های شبکه اجتناب‌ناپذیرند. ممکن است یک پیام ارسال شود، اما تأیید آن گم شود و منجر به تلاش مجدد توسط تولیدکننده گردد. بدون محافظت‌ها، این امر منجر به ایجاد تکراری‌ها می‌شود. در حالی که برخی سیستم‌ها می‌توانند از طریق حذف تکراری‌ها در سمت مصرف‌کننده این موضوع را تحمل کنند، سیستم‌های دیگر نیاز به تضمین سخت‌گیرانه دارند که هر رکورد دقیقاً یک‌بار پردازش شود. کافکا این موضوع را در سطح تولیدکننده از طریق تولیدکننده‌های تکرارپذیر (Idempotent) و در سطح سر تا سر از طریق تراکنش‌ها حل می‌کند.

پیاده‌سازی تولیدکننده‌های تکرارپذیر

تکرارپذیری تضمین می‌کند که یک رکورد یکسان چندین بار در یک پارتیشن واحد به یک موضوع (Topic) نوشته نشود. این امر توسط پیکربندی تولیدکننده enable.idempotence=true کنترل می‌شود. هنگامی که فعال است، تولیدکننده یک شماره ترتیب برای هر پارتیشن حفظ می‌کند. بروتکر این شماره‌های ترتیب را ردیابی می‌کند و پیام‌های خارج از ترتیب یا تکراری از همان تولیدکننده را رد می‌کند.

Properties props = new Properties();
props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
// Enable idempotent producer
props.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, true);

KafkaProducer<String, String> producer = new KafkaProducer<>(props);
producer.send(new ProducerRecord<String, String>("my-topic", "key", "value"));

نکته مهم این است که تولیدکننده‌های تکرارپذیر فقط حذف تکراری‌ها را در یک پارتیشن واحد تضمین می‌کنند. اگر پیامی تلاش مجدد شود و به پارتیشن دیگری مسیریابی شود (به دلیل تغییر کلید یا مشکلات پارتیشنیور)، تکرارپذیری نمی‌تواند از تکراری‌ها بین پارتیشن‌ها جلوگیری کند.

مقیاس‌پذیری به سمت معنای دقیقاً یک‌بار سر تا سر

برای دستیابی به معنای واقعی دقیقاً یک‌بار در چندین موضوع و سیستم‌های خارجی، کافکا API تراکنشی را ارائه می‌دهد. این امکان به تولیدکننده‌ها می‌دهد که به صورت اتمیک به چندین موضوع بنویسند. مصرف‌کننده‌ها می‌توانند با استفاده از تنظیم isolation.level=read_committed در این تراکنش‌ها شرکت کنند و تضمین می‌شود که فقط داده‌های تراکنشی تأیید شده را می‌خوانند.

String transactionalId = "unique-transaction-id";
producer.initTransactions();
producer.beginTransaction();

try {
    producer.send(new ProducerRecord<String, String>("topic-a", "value1"));
    producer.send(new ProducerRecord<String, String>("topic-b", "value2"));
    producer.commitTransaction();
} catch (KafkaException e) {
    producer.abortTransaction();
}

مدیریت ترتیب و تکراری‌ها در مصرف‌کننده‌ها

حتی با نوشتن دقیقاً یک‌بار، منطق مصرف‌کننده باید تکرارپذیر باشد تا موارد حاشیه‌ای را مدیریت کند که در آن‌ها آفست‌ها ممکن است ریست شوند یا یک بازتوزیع گروه مصرف‌کننده رخ دهد. یک الگوی محکم شامل حفظ یک انبار وضعیت محلی از شناسه‌های پیام‌های پردازش شده یا استفاده از یک پایگاه داده با محدودیت‌های یکتا برای رد تکراری‌ها قبل از پردازش است.

نتیجه‌گیری

دستیابی به معنای دقیقاً یک‌بار در آپاچی کافکا نیازمند یک رویکرد چندلایه است. با تولیدکننده‌های تکرارپذیر برای مدیریت تلاش‌های مجدد پایه شروع کنید، از تراکنش‌های کافکا برای نوشتن اتمیک بین موضوع‌ها استفاده کنید و منطق مصرف‌کننده خود را طوری طراحی کنید که ذاتاً تکرارپذیر باشد. با ترکیب این استراتژی‌ها، می‌توانید جریان‌های رویداد محکم و با پهنای باند بالا بسازید که حتی در برابر خرابی‌های گذرا، یکپارچگی داده را حفظ می‌کنند.

Share: