نکته شماره 4: به ابزارهای خط فرمان تسلط دهید - تولید کننده کنسول کافکا
- مصرف کننده کنسول کافکا
- گودال
- سوابق را حذف کنید
- اضافه کردن هدرها به Kafka Records
- بازیابی هدرها
نکته شماره 1: تحویل پیام و ضمانت دوام را درک کنید
برای دوام داده ، Kafkaproducer دارای تنظیمات پیکربندی ACK است. پیکربندی ACKS مشخص می کند که چه تعداد از تصدیق تولید کننده دریافت می کند تا سابقه تحویل شده به کارگزار را در نظر بگیرد. گزینه های انتخاب شده عبارتند از:
- هیچکدام: تولید کننده سوابق تحویل داده شده را با موفقیت در هنگام ارسال سوابق به کارگزار در نظر می گیرد. این اساساً "آتش و فراموش" است.
- یکی: تهیه کننده منتظر است تا کارگزار سرب تصدیق کند که این رکورد را برای ورود به سیستم خود نوشته است.
- همه: تهیه کننده منتظر تصدیق از کارگزار سرب و از کارگزاران زیر است که آنها با موفقیت رکورد را برای سیاهههای مربوطه نوشته اند.
همانطور که مشاهده می کنید ، تجارت در اینجا وجود دارد-و این با طراحی است زیرا برنامه های مختلف نیازهای مختلفی دارند. شما می توانید با احتمال از دست دادن داده ، توان بیشتری را انتخاب کنید ، یا ممکن است ضمانت دوام داده بسیار بالایی را با هزینه توان پایین تر ترجیح دهید.
حالا بیایید یک ثانیه بگیریم تا کمی درباره سناریو ACKS = ALL صحبت کنیم. اگر سوابق خود را با ACK تنظیم کنید که همه آنها را به یک خوشه از سه کارگزار کافکا تنظیم کرده اند ، این بدان معنی است که در شرایط ایده آل ، کافکا شامل سه ماکت از داده های شما است - یکی برای کارگزار سرب و یکی برای هر دو دنبال کننده. هنگامی که سیاهههای مربوط به هر یک از این ماکت ها همه جبران رکورد یکسانی دارند ، در همگام سازی در نظر گرفته می شوند. به عبارت دیگر ، این ماکت های درون همزمان برای یک پارتیشن موضوع خاص ، همان محتوا را دارند. به تصویر زیر نگاهی بیندازید تا به وضوح تصور کنید که چه خبر است:

اما برخی از ظرافت ها برای استفاده از ACKS = ALL پیکربندی وجود دارد. آنچه مشخص نمی کند این است که چه تعداد از ماکت ها باید همگام باشند. کارگزار سرب همیشه با خودش همگام خواهد بود. اما شما می توانید شرایطی داشته باشید که این دو کارگزار به دلیل پارتیشن های شبکه ، بار ضبط و غیره نتوانند از این کار خودداری کنند. اگر این دو دنبال کننده همگام نباشند ، تولید کننده هنوز تعداد مورد نیاز ACK را دریافت می کند ، اما این تنها رهبر در این مورد است. مثلا:

با تنظیم ACK = ALL ، شما حق بیمه را در مورد دوام داده های خود قرار می دهید. بنابراین اگر ماکت ها به کار خود ادامه نمی دهند ، به این دلیل است که شما می خواهید استثناء برای سوابق جدید تا زمانی که ماکت ها گرفتار نشوند.
به طور خلاصه ، داشتن تنها یک ماکت در همگام سازی "نامه قانون" را دنبال می کند اما نه "روح قانون". آنچه ما نیاز داریم ضمانت در هنگام استفاده از تنظیمات ACKS = ALL است. ارسال موفق حداقل اکثریت کارگزاران موجود در همگام سازی را شامل می شود.
به همین ترتیب اتفاق می افتد که یک پیکربندی چنین است: min. insync. replicas. پیکربندی Min. Insync. Replicas تعداد ماکت هایی را که باید در همگام سازی برای نوشتن باشد ، تقویت می کند. توجه داشته باشید که پیکربندی min. insync. replicas در سطح کارگزار یا موضوع تنظیم شده است و پیکربندی تولید کننده نیست. مقدار پیش فرض برای min. insync. replicas یکی است. بنابراین برای جلوگیری از سناریویی که در بالا توضیح داده شد ، در یک خوشه سه کارگزار ، می خواهید مقدار را به دو افزایش دهید.
بیایید نمونه قبلی خود را از قبل دوباره بررسی کنیم و تفاوت را ببینیم:

اگر تعداد ماکت هایی که در همگام سازی هستند زیر مبلغ پیکربندی شده باشد ، کارگزار سرب سعی نمی کند رکورد را به ورود به سیستم خود اضافه کند. رهبر یا notenoughreplicasexception یا notenoughreplicasafterappendexception را پرتاب می کند و تولید کننده را مجبور می کند تا نوشتن را دوباره امتحان کند. داشتن ماکت های خارج از همگام با رهبر یک خطای قابل تکرار در نظر گرفته می شود ، بنابراین تولید کننده همچنان به آزمایش مجدد و ارسال سوابق تا زمان تحویل پیکربندی شده ادامه خواهد داد.
بنابراین با تنظیم تنظیمات Min. Insync. Replicas و تولید کننده ACK برای همکاری با این روش ، دوام داده های خود را افزایش داده اید.
حال بیایید به موارد بعدی در لیست ما حرکت کنیم: پیشرفت در مشتریان Kafka. طی یک سال گذشته ، تولید کننده کافکا و API های مصرف کننده کافکا برخی از ویژگی های جدید را اضافه کرده اند که هر توسعه دهنده کافکا باید آن را بشناسد.
نکته شماره 2: در مورد پارتیشن چسبنده جدید در API تولید کننده اطلاعات کسب کنید
کافکا از پارتیشن ها برای افزایش توان و پخش بار پیام ها به همه کارگزاران در یک خوشه استفاده می کند. سوابق Kafka در قالب کلید/ارزش قرار دارند ، جایی که کلیدها می توانند تهی باشند. تولیدکنندگان کافکا بلافاصله سوابق ارسال نمی کنند ، در عوض آنها را در دسته های خاص پارتیشن قرار می دهند که بعداً ارسال می شوند. دسته ها وسیله ای مؤثر برای افزایش استفاده از شبکه هستند. سه روش وجود دارد که پارتیشن مشخص می کند که در کدام پارتیشن سوابق باید نوشته شود.
پارتیشن را می توان به صراحت در شیء تولید کننده از طریق سازنده اضافه بار تولید کننده ارائه داد. در این حالت ، تولید کننده همیشه از این پارتیشن استفاده می کند.
اگر هیچ پارتیشن ارائه نشده باشد ، و ProductRecord یک کلید دارد ، تولید کننده هش از مدول اصلی را تعداد پارتیشن ها می گیرد. تعداد حاصل از آن محاسبه پارتیشنی است که تولید کننده از آن استفاده می کند.
اگر هیچ کلید و پارتیشن موجود در ProductRecord وجود نداشته باشد ، قبلاً کافکا از یک رویکرد دور رابین برای اختصاص پیام ها در پارتیشن ها استفاده می کرد. تهیه کننده رکورد اول را در گروه به پارتیشن صفر اختصاص می دهد ، دوم به پارتیشن یک و غیره ، تا پایان پارتیشن ها. سپس تولید کننده با پارتیشن صفر شروع می شود و کل فرآیند را برای همه سوابق باقی مانده تکرار می کند.
تصویر زیر این روند را نشان می دهد:

رویکرد دور رابین برای توزیع حتی سوابق موجود در پارتیشن ها به خوبی کار می کند. اما یک اشکال وجود داردبا توجه به این رویکرد دور رابین دور "منصفانه" ، می توانید در نهایت به ارسال دسته های چند جمعیتی کم جمعیت بپردازید. ارسال دسته های کمتری با سوابق بیشتر در هر دسته کارآمدتر است. دسته های کمتری به معنای صف کمتر درخواست تولید هستند ، از این رو بار کمتری روی کارگزاران می شود.
بیایید به یک مثال ساده نگاه کنیم که در آن موضوعی با سه پارتیشن برای توضیح این موضوع دارید. به خاطر سادگی ، فرض کنیم که برنامه شما نه رکورد بدون کلید تولید کرده است ، همه در همان زمان وارد می شوند:

همانطور که در بالا مشاهده می کنید ، 9 رکورد ورودی منجر به سه دسته از سه رکورد خواهد شد. اما بهتر است اگر بتوانیم یک دسته از نه رکورد را ارسال کنیم. همانطور که قبلاً گفته شد ، تعداد کمتری به ترافیک شبکه کمتر و بار کمتری بر روی کارگزاران منجر می شود.
Apache Kafka 2. 4. 0 رویکرد پارتیشن چسبنده را اضافه کرد ، که اکنون این کار را امکان پذیر می کند. به جای استفاده از یک رویکرد دور رابین در هر رکورد ، Partitioner Sticky سوابق را به همان پارتیشن اختصاص می دهد تا اینکه دسته ارسال شود. سپس ، پس از ارسال یک دسته ، پارتیشن چسبنده پارتیشن را افزایش می دهد تا برای دسته بعدی استفاده شود. بیایید تصویر خود را از بالا مجدداً مورد استفاده قرار دهیم اما با استفاده از قسمت چسبنده چسبنده به روز شده ایم:

با استفاده از همان پارتیشن تا زمانی که یک دسته پر شود یا در غیر این صورت تکمیل شود ، ما درخواست های تولید کمتری را ارسال خواهیم کرد ، که این باعث کاهش بار در صف درخواست می شود و تأخیر سیستم را نیز کاهش می دهد. شایان ذکر است که پارتیشل چسبنده هنوز منجر به توزیع یکنواخت سوابق می شود. توزیع یکنواخت با گذشت زمان اتفاق می افتد ، زیرا قسمت تقسیم کننده برای هر پارتیشن دسته ای می فرستد. شما می توانید آن را به عنوان یک رویکرد دور رابین یا "در نهایت یکنواخت" دور کنید.
برای کسب اطلاعات بیشتر در مورد Partitioner Sticky ، می توانید پیشرفت های تولید کننده Apache Kafka را با پست وبلاگ Sticky Partitioner و سند طراحی KIP-480 مرتبط بخوانید.
اکنون بیایید به تغییرات مصرف کننده حرکت کنیم.
نکته شماره 3: با استفاده از مجدد تعادل تعاونی ، گروه های مصرف کننده "Stop-the-World" را خودداری کنید
کافکا یک سیستم توزیع شده است و یکی از مهمترین کارهایی که سیستم های توزیع شده باید انجام دهند ، مقابله با خرابی ها و اختلالات است - نه فقط پیش بینی خرابی ها ، بلکه کاملاً آنها را در آغوش بگیرید. یک نمونه عالی از نحوه برخورد کافکا این اختلال مورد انتظار ، پروتکل گروه مصرف کننده است که چندین نمونه از مصرف کننده را برای یک برنامه منطقی واحد مدیریت می کند. اگر نمونه ای از مصرف کننده متوقف شود ، با طراحی یا در غیر این صورت ، کافکا دوباره تعادل می یابد و مطمئن می شود که نمونه دیگری از مصرف کننده کار را به دست می گیرد.
از نسخه 2. 4 ، کافکا یک پروتکل جدید تعادل را معرفی کرد: تعادل تعاونی. اما قبل از اینکه به پروتکل جدید شیرجه بزنیم ، بیایید جزئیات کمی بیشتر به اصول گروه مصرف کننده نگاه کنیم.
بیایید فرض کنیم شما یک برنامه توزیع شده با چندین مصرف کننده مشترک در یک موضوع دارید. هر مجموعه ای از مصرف کنندگان که با همان گروه پیکربندی شده اند. یک مصرف کننده منطقی به نام یک گروه مصرف کننده را تشکیل می دهد. هر مصرف کننده در این گروه وظیفه مصرف از یک یا چند پارتیشن از موضوع (های) مشترک را بر عهده دارد. این پارتیشن ها توسط رهبر گروه مصرف کننده اختصاص یافته است.
در اینجا تصویری که این مفهوم را نشان می دهد وجود دارد:

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

همانطور که مشاهده می کنید ، مصرف کننده 2 به دلایلی ناکام است و یا نظرسنجی را از دست می دهد یا باعث ایجاد زمان جلسه می شود. هماهنگ کننده گروه آن را از گروه خارج می کند و آنچه را که به عنوان یک تعادل شناخته می شود ، محرک می کند. تعادل مکانیسمی است که سعی در توزیع یکنواخت (تعادل) بار کار در تمام اعضای موجود در یک گروه مصرف کننده دارد. در این حالت ، از آنجا که مصرف کننده 2 گروه را ترک کرد ، Reflace پارتیشن های قبلی خود را به سایر اعضای فعال گروه اختصاص می دهد. بنابراین همانطور که می بینید ، از دست دادن یک برنامه مصرف کننده برای یک شناسه گروه خاص منجر به از بین رفتن پردازش در آن پارتیشن های موضوعی نمی شود.
با این حال ، اشکال از رویکرد پیش فرض تعادل وجود دارد. هر مصرف کننده تمام تکلیف خود را از پارتیشن های موضوع واگذار می کند ، و تا زمانی که پارتیشن های موضوع مجدداً تنظیم شوند ، هیچ پردازشی صورت نمی گیرد-گاهی اوقات به عنوان یک تعادل "توقف جهان" گفته می شود. برای ترغیب این مسئله ، بسته به نمونه ConsumerPartitionAssignor مورد استفاده ، مصرف کنندگان به سادگی با همان پارتیشن های موضوعی که قبل از تعادل متعلق به آنها بودند ، مجدداً مجدداً انتصاب می کنند ، تأثیر خالص این است که نیازی به مکث کار در آن پارتیشن ها نیست.
این اجرای پروتکل rebalance به نام Reeger تعادل نامیده می شود زیرا اهمیت اطمینان از این که هیچ مصرف کننده در همان گروه مالکیت خود را بر روی همان پارتیشن های موضوعی ادعا نمی کند ، در اولویت قرار می دهد. مالکیت همان پارتیشن موضوع توسط دو مصرف کننده در همان گروه منجر به رفتار نامشخص می شود.
در حالی که مهم است که هر دو مصرف کننده از ادعای مالکیت در همان پارتیشن موضوع جلوگیری کنند ، اما معلوم می شود که رویکرد بهتری وجود دارد که ایمنی را بدون به خطر انداختن در زمان صرف پردازش نمی کند: تعادل تعاونی افزایشی. برای اولین بار به Kafka Coect در Apache Kafka 2. 3 معرفی شد ، این هم اکنون برای پروتکل گروه مصرف کننده نیز اجرا شده است. با رویکرد تعاونی ، مصرف کنندگان به طور خودکار از مالکیت همه پارتیشن های موضوع در ابتدای تعادل خودداری نمی کنند. در عوض ، همه اعضا وظیفه فعلی خود را رمزگذاری می کنند و اطلاعات را به رهبر گروه ارسال می کنند. رهبر گروه سپس تعیین می کند که کدام پارتیشن ها نیاز به تغییر مالکیت دارند - در نظر گرفتن تولید یک تکلیف کاملاً جدید از ابتدا.
اکنون یک تعادل دوم صادر می شود ، اما این بار ، فقط پارتیشن های موضوعی که نیاز به تغییر مالکیت دارند ، درگیر هستند. این می تواند پارتیشن های موضوعی را که دیگر به آنها اختصاص داده نشده یا اضافه کردن پارتیشن های موضوع جدید است ، ابطال کند. برای پارتیشن های موضوعی که در هر دو تکلیف جدید و قدیمی قرار دارند ، هیچ چیز نباید تغییر کند ، این بدان معنی است که پردازش ادامه برای پارتیشن های موضوعی که در حال حرکت نیستند.
نکته آخر این است که از بین بردن رویکرد "متوقف کردن جهان" برای تعادل مجدد و فقط متوقف کردن پارتیشن های موضوع درگیر به معنای تغییر مجدد کم هزینه است ، بنابراین کل زمان را برای تکمیل تعادل کاهش می دهد. حتی در حال حاضر دوباره تعادل طولانی دردناک تر است که پردازش می تواند در سراسر آنها ادامه یابد. این تغییر مثبت در تعادل با استفاده از COOPERATIVESTICKYASSINOR امکان پذیر است. CooperativeStickyassignor تجارت با تعادل دوم را انجام می دهد اما با بهره گیری از بازگشت سریعتر به عملیات عادی.
برای فعال کردن این پروتکل جدید regalance ، شما باید برای استفاده از COOPERATIVESTICKYASSINOR ، Partition. Assignment. Strategy را تنظیم کنید. همچنین ، توجه داشته باشید که این تغییر کاملاً در سمت مشتری است. برای استفاده از پروتکل جدید Rebalance ، فقط باید نسخه مشتری خود را به روز کنید. اگر شما یک کاربر Kafka Streams هستید ، خبرهای بهتری نیز وجود دارد. Kafka Streams به طور پیش فرض پروتکل تعادل تعاونی را امکان پذیر می کند ، بنابراین کار دیگری وجود ندارد.
نکته شماره 4: به ابزارهای خط فرمان تسلط دهید
نصب باینری Apache Kafka شامل چندین ابزار واقع در فهرست سطل است. در حالی که شما چندین ابزار در آن فهرست پیدا خواهید کرد ، می خواهم چهار ابزار را به شما نشان دهم که فکر می کنم بیشترین تأثیر را در کارهای روزانه شما داشته باشد. من به کنسول-مصرف کننده ، کنسول-تولید کننده ، دامپزشک و حذف رکورد مراجعه می کنم.
تولید کننده کنسول کافکا
تولید کننده کنسول به شما امکان می دهد سوابق مربوط به یک موضوع را مستقیماً از خط فرمان تهیه کنید. تولید از خط فرمان راهی عالی برای آزمایش سریع برنامه های جدید مصرف کننده در هنگام تولید اطلاعات به موضوعات نیست. برای شروع تولید کننده کنسول ، این دستور را اجرا کنید:
kafka-console-تولید کننده--موضوعی -لیست کارگزاران بعد از اجرای دستور ، یک سریع خالی در انتظار ورودی شما وجود دارد - فقط برخی از کاراکترها را تایپ کرده و برای تولید یک پیام به Enter ضربه بزنید.
استفاده از تولید کننده خط فرمان از این طریق هیچ کلید و فقط مقادیر ارسال نمی کند. خوشبختانه ، راهی برای ارسال کلیدها نیز وجود دارد. شما فقط باید دستور را به روز کنید تا پرچم های لازم را در بر بگیرید:
kafka-console-producer--top -broker-list -property parse. key = true -property key. separator = ":"
انتخاب ویژگی key. peparator دلخواه است. می توانید از هر کاراکتر استفاده کنید. و اکنون ، می توانید جفت کلید/مقدار کامل را از خط فرمان ارسال کنید! اگر از رجیستری طرحواره Confluent استفاده می کنید ، تولید کنندگان خط فرمان برای ارسال سوابق در قالب های Avro ، ProtoBUF و JSON Schema در دسترس هستند.
حال بیایید نگاهی به طرف دیگر سکه بیندازیم: سوابق از خط فرمان.
مصرف کننده کنسول کافکا
مصرف کننده کنسول به شما امکان مصرف سوابق از یک موضوع Kafka را مستقیماً از خط فرمان می دهد. قادر به شروع سریع مصرف کننده می تواند ابزاری ارزشمند در نمونه سازی یا اشکال زدایی باشد. در نظر بگیرید که یک میکروسرویس جدید ایجاد کنید. برای تأیید سریع که برنامه تولید کننده شما در حال ارسال پیام است ، می توانید به سادگی این دستور را اجرا کنید:
kafka-console-consumer--top -bootstrap-server بعد از اجرای این دستور ، شروع به دیدن سوابق در صفحه نمایش خود خواهید کرد (تا زمانی که داده ها در حال حاضر به موضوع تولید می شوند). اگر می خواهید تمام سوابق را از ابتدا مشاهده کنید ، می توانید یک پرچم-از کارگردانی را به دستور اضافه کنید ، و تمام سوابق تولید شده در آن موضوع را مشاهده خواهید کرد.
kafka-console-consumer--top -bootstrap-server -از بین بردن
اگر از Registry Schema استفاده می کنید ، مصرف کنندگان خط فرمان برای سوابق رمزگذاری شده AVRO ، ProTOBUF و JSON در دسترس هستند. مصرف کنندگان خط فرمان رجیستری Schema برای کار با سوابق در قالب های AVRO ، ProTOBUF یا JSON در نظر گرفته شده اند ، در حالی که مصرف کنندگان ساده با سوابق نوع جاوا ابتدایی کار می کنند: رشته ، طولانی ، دوتایی ، عدد صحیح و غیره. قالب پیش فرض مورد انتظار برای کلیدها و کلیدهامقادیر توسط مصرف کننده ساده کنسول نوع رشته است.
اگر کلیدها یا مقادیر رشته ای نیستند ، باید از طریق پرچم های خط فرمان-کلید-دسیریالایزر و-Value-deserializer با نام کلاس کاملاً واجد شرایط از deserializer های مربوطه ، deserializer را تهیه کنید.
شاید به خوبی متوجه شده باشید که به طور پیش فرض ، مصرف کننده کنسول فقط مؤلفه ارزش پیام ها را به صفحه چاپ می کند. اگر می خواهید کلیدها را نیز ببینید ، می توانید این کار را با درج پرچم های لازم انجام دهید:
kafka-console-consumer--top -bootstrap-server -property print. key = true-property key. separator = ":"
مانند تولید کننده ، مقدار مورد استفاده برای جداکننده کلید دلخواه است ، بنابراین می توانید هر شخصیتی را که می خواهید استفاده کنید انتخاب کنید.
گودال
بعضی اوقات وقتی با کافکا کار می کنید ، ممکن است خود را پیدا کنید که نیاز به بازرسی دستی در زمینه های اساسی یک موضوع داشته باشید. این که آیا شما فقط در مورد داخلی کافکا کنجکاو هستید یا باید یک مسئله را اشکال زدایی کنید و محتوا را تأیید کنید ، دستور kafka-dump-log دوست شماست. در اینجا یک دستور استفاده شده برای مشاهده ورود به یک موضوع به عنوان مثال به عنوان مثال:
kafka-dump-log -print-data-log -files ./var/lib/kafka/data/example-0/000000000000000000. log
- پرچم-print-data-log برای چاپ داده ها در سیاهه مشخص می کند.
- پرچ م-پرونده مورد نیاز است. این همچنین می تواند یک لیست جدا از کاما باشد.
برای یک لیست کامل از گزینه ها و توضیحات مربوط به هر گزینه ، Kafka-Dump-Log را با پرچم-Help اجرا کنید.
اجرای دستور بالا چیزی شبیه به این است:
دامپینگ ./var/lib/kafka/data/example-0/000000000000000000. LOG شروع شروع: 0 پایه پایه: 0 LastOffset: 0 تعداد: 1 پایه: -1 آخرین توالی: -1 تولید کننده: -1 ProciationPoch: -1 PartitionLeaderePoch: 0 0iStransactional: false iscontrol: موقعیت کاذب: 0 CreateTime: 1599775774460 اندازه: 81 جادو: 2 CompressCodec: None CRC: 3162584294 ISVALID: TRUE |افست: 0 CreateTime: 1599775774460 Keysize: 3 مقادیر: 10 دنباله: -1 هدرکی ها: [] کلید: 887 بار: -2. 1510235 پایه: 1 LastOffset: 9 Count: 9 Basesequence: -1 آخرین نسخه: -1 تولید کننده: -1 Producterepoch:-1 PartitionLeaderePoch: 0 iStransactional: False ISControl: False Position: 81 CreateTime: 159977574468 Size: 252 Magic: 2 CompressCodec: None CRC: 2796351311 isValid: True |افست: 1 CreateTime: 1599775774463 Keysize: 1 مقادیر: 9 دنباله: -1 هدر: [] کلید: 5 بار بار: 33. 440664 |افست: 2 CreateTime: 1599775774463 Keysize: 8 مقادیر: 9 دنباله: -1 HeaderKeys: [] کلید: 60024247 بار: 9. 1408728 |افست: 3 CreateTime: 1599775774463 Keysize: 1 مقادیر: 9 دنباله: -1 هدر: [] کلید: 1 بار بار: 45. 348946 |افست: 4 CreateTime: 1599775774464 Keysize: 6 مقادیر: 10 دنباله: -1 هدر: [] کلید: 241795 بار: -63. 786373 |افست: 5 CreateTime: 1599775774465 Keysize: 8 مقادیر: 9 دنباله: -1 HeaderKeys: [] کلید: 53596698 بار: 69. 431393 |افست: 6 CreateTime: 1599775774465 Keysize: 8 مقادیر: 9 دنباله: -1 HeaderKeys: [] کلید: 33219463 بار: 88. 307875 |افست: 7 CreateTime: 1599775774466 Keysize: 1 مقادیر: 9 دنباله: -1 هدر: [] کلید: 0 بار بار: 39. 940350 |افست: 8 CreateTime: 1599775774467 Keysize: 5 مقادیر: 9 دنباله: -1 HeaderKeys: [] کلید: 78496 Payload: 74. 180098 |افست: 9 CreateTime: 1599775774468 Keysize: 8 مقادیر: 9 دنباله: -1 هدر: [] کلید: 89866187 بار: 79. 459314
اطلاعات زیادی از دستور dump-log در دسترس است. برای هر رکورد می توانید کلید ، بار بار (مقدار) ، جبران و جدول زمانی را به وضوح مشاهده کنید. به خاطر داشته باشید که این داده ها از یک موضوع نمایشی است که فقط شامل 10 پیام است ، بنابراین با یک موضوع واقعی ، داده های قابل توجهی بیشتر خواهد بود. همچنین توجه داشته باشید که در این مثال ، کلیدها و مقادیر موضوع رشته ها هستند. برای اجرای ابزار dump-log با انواع کلید یا ارزش به غیر از رشته ها ، باید از پرچم های کلاس-key-decoder یا-value-decoder-class استفاده کنید.
سوابق را حذف کنید
Kafka سوابق مربوط به موضوعات مربوط به دیسک را ذخیره می کند و حتی هنگامی که مصرف کنندگان آن را خوانده اند ، داده ها را حفظ می کند. با این حال ، سوابق در یک پرونده بزرگ ذخیره نمی شوند ، اما با پارتیشن در بخش هایی که ترتیب جبران در بخش های مختلف برای یک پارتیشن موضوعی مداوم است ، به بخش هایی تقسیم نمی شوند. از آنجا که سرورها مقدار نامحدودی از ذخیره سازی ندارند ، کافکا تنظیماتی را برای کنترل میزان حفظ داده ها بر اساس زمان و اندازه فراهم می کند:
- تنظیمات زمان کنترل کنترل داده ها log. retent. hours است که به طور پیش فرض به 168 ساعت (یک هفته)
- اندازه پیکربندی اندازه log. retent. Bytes کنترل می کند که بخش های بزرگ قبل از واجد شرایط حذف می توانند رشد کنند
با این حال ، تنظیم پیش فرض برای log. retention. byte s-1 است که اجازه می دهد اندازه بخش ورود به سیستم نامحدود باشد. اگر مراقب نیستید و اندازه نگهداری و همچنین زمان نگهداری را پیکربندی نکرده اید ، می توانید وضعیتی داشته باشید که فضای دیسک را تمام کنید. به یاد داشته باشید ، شما هرگز نمی خواهید وارد سیستم فایل شوید و پرونده ها را به صورت دستی حذف کنید. در عوض ، ما می خواهیم یک روش کنترل شده و پشتیبانی شده برای حذف سوابق از یک موضوع به منظور آزاد کردن فضا. خوشبختانه ، کافکا با ابزاری که داده های لازم را حذف می کند ، ارسال می شود.
رکوردهای kafka-delete دارای دو پارامتر اجباری است:
- -bootstrap-server: کارگزار (های) برای اتصال به bootstrapping
- -Offset-Json-File: یک فایل JSON که حاوی تنظیمات حذف است
در اینجا نمونه ای از پرونده JSON آورده شده است:
<"partitions": [ ], "version":1>
همانطور که مشاهده می کنید ، قالب JSON ساده است. این مجموعه ای از اشیاء JSON است. هر شیء JSON سه ویژگی دارد:
- موضوع: موضوع برای حذف از
- پارتیشن: پارتیشن برای حذف
- افست: افست که می خواهید حذف از آن شروع شود ، به سمت عقب حرکت کنید تا جبران شود
برای این مثال ، من در حال استفاده مجدد از همان موضوع از ابزار dump-log هستم ، بنابراین این یک فایل JSON بسیار ساده است. اگر پارتیشن ها یا موضوعات بیشتری داشتید ، به سادگی در پرونده پیکربندی JSON در بالا گسترش می دهید.
من می خواهم در مورد نحوه انتخاب افست در پرونده پیکربندی JSON بحث کنم. از آنجا که موضوع مثال فقط شامل 10 سوابق است ، می توانید به راحتی جبران شروع را برای شروع فرآیند حذف محاسبه کنید. اما در عمل ، شما به احتمال زیاد نمی دانید از بالای سر خود استفاده کنید که چه چیزی را جبران کنید. همچنین به خاطر داشته باشید که جبران! = شماره پیام ، بنابراین نمی توانید فقط از "پیام 42" حذف کنیداگر ی ک-1 را تهیه کنید ، از جبران Watermark بالا استفاده می شود ، به این معنی که تمام داده های موجود در این موضوع را حذف خواهید کرد. Watermark High بالاترین جبران در دسترس برای مصرف است (جبران آخرین پیام با موفقیت تکرار شده ، به علاوه یک).
اکنون برای اجرای دستور ، فقط این کار را در خط فرمان وارد کنید:
kafka-delete-Records-BootStrap-Server -Offset-Json-File Offsets. json
پس از اجرای این دستور ، باید چیزی شبیه به این را در کنسول ببینید:
اجرای سوابق حذف سوابق عملکرد حذف عملیات تکمیل شده: پارتیشن: مثا ل-0 Low_watermark: 10
نتایج این فرمان نشان می دهد که کافکا تمام سوابق را از نمونه موضوع نمون ه-0 حذف کرده است. مقدار Low_watermark 10 نشان دهنده کمترین جبران در دسترس مصرف کنندگان است. از آنجا که فقط 10 سوابق در موضوع مثال وجود داشت ، می دانیم که جبران خسارات از 0 تا 9 و هیچ مصرف کننده نمی تواند دوباره آن سوابق را بخواند. برای پیش زمینه بیشتر در مورد اجرای حذف ، می توانید KIP-107 و KIP-204 را بخوانید.
نکته شماره 5: از قدرت هدرهای ضبط استفاده کنید
Apache Kafka 0. 11 مفهوم هدرهای رکورد را معرفی کرد. هدرهای ضبط به شما امکان اضافه کردن برخی از ابرداده در مورد رکورد Kafka را می دهند ، بدون اینکه اطلاعات اضافی به جفت کلید/ارزش خود رکورد اضافه کنید. در نظر بگیرید که آیا می خواهید برخی از اطلاعات را در یک پیام ، مانند شناسه برای سیستمی که داده ها از آن سرچشمه گرفته است ، جاسازی کنید. شاید شما این را برای اهداف سلسله و حسابرسی و به منظور تسهیل مسیریابی داده های پایین دست بخواهید.
چرا فقط این اطلاعات را به کلید اضافه نمی کنید؟سپس می توانید قسمت مورد نیاز را استخراج کنید و می توانید مطابق آن داده ها را مسیریابی کنید. اما افزودن داده های مصنوعی به کلیدی دو مشکل بالقوه را ایجاد می کند.
- اول ، اگر از یک موضوع فشرده استفاده می کنید ، اضافه کردن اطلاعات به کلید باعث می شود که ضبط نادرست به صورت منحصر به فرد ظاهر شود. بنابراین ، تراکم مطابق آنچه در نظر گرفته شده عمل نمی کند.
- برای شماره دوم ، اگر یک شناسه خاص سیستم در سوابق ارسال شده حاکم باشد ، تأثیر را در نظر بگیرید. شما اکنون شرایطی دارید که می توانید یک کلید مهم را داشته باشید. بسته به نحوه مصرف شما از پارتیشن ها ، توزیع ناهموار کلیدها می تواند با افزایش تأخیر در پردازش تأثیر بگذارد.
این دو موقعیت است که ممکن است بخواهید از هدرها استفاده کنید. عنوان اصلی KIP ارائه شده موارد دیگری را نیز ارائه می دهد:
- مسیریابی خودکار پیام ها بر اساس اطلاعات هدر بین خوشه ها
- ابزارهای APM سازمانی (به عنوان مثال ، AppDynamics یا Dynatrace) باید در شناسه های معامله "جادویی" بخیه برای آنها برای ارائه نظارت بر جریان معاملات پایان به پایان.
- ابرداده حسابرسی با پیام ضبط می شود ، به عنوان مثال ، مشتری-ID که رکورد را تولید می کند.
- برای پایان دادن به پایان کار باید به پایان برسد و بدون دستکاری امضا شود ، اما اجزای اکوسیستم برای دستیابی به وظایف نیاز به دسترسی به ابرداده دارند.
اکنون که من برای استفاده از هدرها موردی ایجاد کرده ام ، اجازه دهید چگونه بتوانید سرصفحه ای را به سوابق Kafka خود اضافه کنید.
اضافه کردن هدرها به Kafka Records
در اینجا کد جاوا برای اضافه کردن عنوان به یک تولید کننده ارائه شده است:
ProductRecord ProductRecord = New ProductRecord<>("bizops" ، "مقدار") ؛ ProductRecord. Headers (). افزودن ("مشتری-ID" ، "2334" . getBytes (StandardCharsets. utf_8)) ؛ProductRecord. Headers (). افزودن ("پرونده داده" ، "ورودی-data. txt" . getBytes (StandardCharsets. utf_8)) ؛
// جزئیات باقی مانده برای Clarity Builder. Send (ProductRecord) ؛
- نمونه ای از کلاس Producrecord ایجاد کنید
- با روش Processerrecord. Headers () تماس بگیرید و کلید و مقدار را برای هدر اضافه کنید
- اضافه کردن یک هدر دیگر
چند مورد وجود دارد که باید در اینجا با مثال کد به آن اشاره کنیم. رابط هدر انتظار یک کلید رشته و مقدار به عنوان یک آرایه بایت را دارد. حتی اگر یک کلید را ارائه می دهید ، در صورت لزوم می توانید به همان تعداد هدر با همان کلید اضافه کنید. کلیدهای تکراری ورودی های قبلی را با همان کلید بازنویسی نمی کنند.
همچنین ، سازندگان ProdureCord بیش از حد بارگذاری شده وجود دارند که یک تکرار پذیر را می پذیرند. شما می توانید کلاس بتونی خود را ایجاد کنید که رابط هدر را پیاده سازی می کند و در مجموعه ای که رابط قابل تکرار را پیاده سازی می کند ، عبور می کند. با این حال ، در عمل ، روش ساده نشان داده شده در اینجا باید کافی باشد.
اکنون که می دانید چگونه هدر را اضافه کنید ، بیایید نگاهی بیندازیم که چگونه می توانید از طرف مصرف کننده چیزها به هدرها دسترسی پیدا کنید.
بازیابی هدرها
اینگونه است که می توانید هنگام مصرف سوابق به عنوان ها دسترسی پیدا کنید:
// جزئیات باقی مانده برای وضوح
ConsumerRecords ConsumerRecords = Consumer. poll (مدت زمان. ofseconds (1)) ؛
برای (ConsumerRecord ConsumerRecord: ConsumerRecords)<for (header="" header="" :="" consumerrecord.headers())="" system.out.println("header="" key="" "="" +="" header.key()="" "header="" value="" new="" string(header.value()));="">>
- تکرار بیش از ConsumerRecords
- برای هر ConsumerRecord ، تکرار بر روی هدرها
- پردازش هدر
از کد فوق ، می بینید که برای پردازش هدرها ، به سادگی از روش ConsumerRecord. Headers () برای بازگشت عنوان ها استفاده کنید. در مثال بالا ، ما برای اهداف تظاهرات ، هدرها را به سمت کنسول چاپ می کنیم. پس از دسترسی به هدرها ، می توانید آنها را در صورت لزوم پردازش کنید. برای خواندن هدرهای از خط فرمان ، KIP-431 پشتیبانی از هدرهای چاپی اختیاری را از Consoleconsumer اضافه می کند ، که در نسخه Apache Kafka 2. 7. 0 در دسترس خواهد بود.
همچنین می توانید از KAFKACAT برای مشاهده هدرها از خط فرمان استفاده کنید. در اینجا یک دستور مثال وجود دارد:
kafkaca t-b kafk a-broker: 909 2-t my_topic_nam e-c -f ' nkey (٪ k bytes): ٪ k مقدار (٪ s bytes): ٪ s timestamp: ٪ t پارتیشن: ٪ p جبران: ٪ o عنوان ها:٪ H n '
درج کردن
اکنون پنج نکته برتر برای همکاری با Apache Kafka را خوانده اید. برای یادآوری ، ما می فهمیم:
- دوام پیام و رابطه آن با ضمانت های تحویل
- پارتیشن چسبنده در API تولید کننده
- ابزارهای خط فرمان
- قدرت هدرهای رکورد
و با این حال ، هنوز چیزهای بیشتری برای یادگیری وجود دارد! برای دیدن آنچه اتفاق می افتد ، به آموزش های توسعه دهنده و Kafka بروید.