مفاهیم جریان

ساخت وبلاگ

در این بخش مفاهیم کلیدی جریان های کافکا را خلاصه می کنیم. برای کسب اطلاعات بیشتر ، لطفاً به راهنمای Streams Architecture و راهنمای توسعه دهنده جریان مراجعه کنید. شما همچنین ممکن است به دوره Kafka Streams 101 علاقه مند باشید.

کافکا 101¶

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

  • Who Who: Kafka تولید کنندگان ، مصرف کنندگان و کارگزاران را متمایز می کند. به طور خلاصه ، تولید کنندگان داده ها را به کارگزاران کافکا منتشر می کنند ، و مصرف کنندگان داده های منتشر شده از کارگزاران کافکا را می خوانند. تولیدکنندگان و مصرف کنندگان کاملاً جدا شده اند و هر دو در خارج از کارگزاران کافکا در محیط یک خوشه کافکا اجرا می شوند. یک خوشه کافکا از یک یا چند کارگزار تشکیل شده است. برنامه ای که از Kafka Streams API استفاده می کند ، هم به عنوان تولید کننده و هم مصرف کننده عمل می کند.
  • داده ها: داده ها در موضوعات ذخیره می شوند. موضوع مهمترین انتزاعی است که توسط کافکا ارائه شده است: این یک دسته یا نام خوراک است که داده ها توسط تولید کنندگان منتشر می شود. هر موضوع در کافکا به یک یا چند پارتیشن تقسیم می شود. داده های پارتیشن کافکا برای ذخیره ، حمل و نقل و تکرار آن. کافکا داده های پارتیشن را برای پردازش آن پخش می کند. در هر دو مورد ، این پارتیشن بندی امکان ارتجاعی ، مقیاس پذیری ، عملکرد بالا و تحمل گسل را فراهم می کند.
  • موازی سازی: پارتیشن های موضوعات کافکا و به ویژه تعداد آنها برای یک موضوع معین ، همچنین عامل اصلی است که موازی بودن کافکا را در رابطه با خواندن و نوشتن داده ها تعیین می کند. به دلیل ادغام محکم با کافکا ، موازی سازی برنامه ای که از API جریان Kafka استفاده می کند ، در درجه اول به موازی بودن کافکا بستگی دارد.

جریان

یک جریان مهمترین انتزاعی است که توسط جریان های کافکا ارائه شده است: این یک مجموعه داده بدون مرز و به طور مداوم به روز می کند ، جایی که به معنای "ناشناخته یا از اندازه نامحدود" است. درست مانند یک موضوع در کافکا ، جریانی در API جریان Kafka از یک یا چند پارتیشن جریان تشکیل شده است.

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

برنامه پردازش جریان

یک برنامه پردازش جریان هر برنامه ای است که از کتابخانه جریان Kafka استفاده می کند. در عمل ، این بدان معنی است که احتمالاً برنامه "شما" است. این ممکن است منطق محاسباتی خود را از طریق یک یا چند توپولوژی پردازنده تعریف کند.

برنامه پردازش جریان شما در داخل یک کارگزار اجرا نمی شود. در عوض ، در یک نمونه JVM جداگانه یا در یک خوشه جداگانه کاملاً اجرا می شود.

../_images/streams-apps-not-ruing-in-brokers.png

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

../_images/scale-out-streams-app.png

توپولوژی پردازنده ¶

یک توپولوژی پردازنده یا به سادگی توپولوژی منطق محاسباتی پردازش داده ها را که باید توسط یک برنامه پردازش جریان انجام شود ، تعریف می کند. توپولوژی نمودار پردازنده های جریان (گره ها) است که توسط جریان ها (لبه ها) به هم وصل می شوند. توسعه دهندگان می توانند توپولوژی را از طریق API پردازنده سطح پایین یا از طریق جریان های Kafka DSL ، که در بالای سابق ساخته می شود ، تعریف کنند.

../_images/streams-concepts-topology.jpg

مستندات معماری با جزئیات بیشتری توپولوژی را توصیف می کند.

پردازنده جریان

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

Kafka Streams دو API را برای تعریف پردازنده های جریان فراهم می کند:

  1. DSL اعلانی و عملکردی API توصیه شده برای اکثر کاربران - و به ویژه برای مبتدیان - زیرا بیشتر موارد استفاده از پردازش داده ها فقط در چند خط از کد DSL قابل بیان است. در اینجا ، شما به طور معمول از عملیات داخلی مانند نقشه و فیلتر استفاده می کنید.
  2. API پردازنده ضروری و سطح پایین ، انعطاف پذیری بیشتری را نسبت به DSL در اختیار شما قرار می دهد اما به هزینه نیاز به کار کد نویسی دستی بیشتر است. در اینجا ، می توانید پردازنده های سفارشی و همچنین ارتباط مستقیم با فروشگاه های دولتی را تعریف و وصل کنید.

پردازش جریان مطبوع

برخی از برنامه های پردازش جریان نیازی به حالت ندارند - آنها بدون تابعیت هستند - به این معنی که پردازش یک پیام از پردازش پیام های دیگر مستقل است. مثالها وقتی فقط نیاز به تبدیل یک پیام به طور همزمان دارید یا پیام ها را بر اساس برخی شرایط فیلتر می کنید.

با این حال ، در عمل ، بیشتر برنامه ها برای کار صحیح به دولت نیاز دارند-آنها وضعیتی دارند-و این حالت باید به روشی تحمل به خطا اداره شود. برنامه شما هر زمان که ، به عنوان مثال ، نیاز به پیوستن ، جمع یا پنجره داده های ورودی آن دارد. Kafka Streams قابلیت های پردازش قدرتمند ، الاستیک ، بسیار مقیاس پذیر و تحمل گسل را در اختیار شما قرار می دهد.

دوگانگی جریانها و جداول

هنگام اجرای موارد استفاده از پردازش جریان در عمل ، به طور معمول به جریان و همچنین پایگاه داده نیاز دارید. یک مورد استفاده به عنوان مثال که در عمل بسیار رایج است ، یک برنامه تجارت الکترونیکی است که با آخرین اطلاعات مشتری از جدول پایگاه داده ، جریان ورودی از معاملات مشتری را غنی می کند. به عبارت دیگر ، جریان ها در همه جا وجود دارد ، اما پایگاه داده ها نیز در همه جا هستند.

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

قبل از اینکه ما در مورد مفاهیمی مانند تجمع در جریان های کافکا بحث کنیم ، ابتدا باید جداول را با جزئیات بیشتری معرفی کنیم و در مورد دوگانگی جدول جریان فوق الذکر صحبت کنیم. در اصل ، این دوگانگی به این معنی است که می توان یک جریان را به عنوان یک جدول مشاهده کرد و یک جدول را می توان به عنوان یک جریان مشاهده کرد.

ما عمداً توضیحات زیر را ساده نگه می داریم و بدین ترتیب بحث در مورد کلیدهای مرکب ، چندست و غیره را صرف می کنیم.

یک شکل ساده از یک جدول مجموعه ای از جفت های با ارزش کلیدی است که به آن نقشه یا آرایه انجمنی نیز گفته می شود. چنین جدول ممکن است به شرح زیر باشد:

../_images/streams-table-duality-01.jpg

دوگانگی جدول جریان ، رابطه نزدیک بین جریانها و جداول را توصیف می کند.

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

بیایید این را با یک مثال نشان دهیم. یک جدول را تصور کنید که تعداد کل صفحه نمایش توسط کاربر را ردیابی کند (ستون اول نمودار زیر). با گذشت زمان ، هر زمان که یک رویداد صفحه جدید پردازش شود ، وضعیت جدول بر این اساس به روز می شود. در اینجا ، دولت بین نقاط مختلف در زمان تغییر می کند - و تجدید نظر های مختلف جدول - می تواند به عنوان یک جریان ChangeLog (ستون دوم) نشان داده شود.

../_images/streams-table-duality-02.jpg

به دلیل دوگانگی جدول جریان ، می توان از همان جریان برای بازسازی جدول اصلی (ستون سوم) استفاده کرد:

../_images/streams-table-duality-03.jpg

به عنوان مثال ، از همان مکانیسم برای تکرار بانکهای اطلاعاتی از طریق تغییر ضبط داده ها (CDC) و ، در جریان های کافکا ، برای تکرار فروشگاه های به اصطلاح حالت خود در دستگاه ها برای تحمل گسل استفاده می شود. دوگانگی جدول جریان مفهوم مهمی برای برنامه های پردازش جریان در عمل است که کافکا آن را به صراحت از طریق انتزاع Kstream و Ktable مدل می کند ، که ما در بخش های بعدی توصیف می کنیم.

Kstream¶

فقط جریان های کافکا DSL مفهوم kstream را دارد.

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

برای نشان دادن ، تصور می کنیم دو سوابق داده زیر به جریان ارسال می شوند:

اگر برنامه پردازش جریان شما برای جمع آوری مقادیر برای هر کاربر باشد ، 4 برای آلیس باز می گردد. چرا؟از آنجا که سوابق داده دوم به روزرسانی رکورد قبلی در نظر گرفته نمی شود. این رفتار Kstream را با Ktable در زیر مقایسه کنید ، که 3 برای آلیس باز می گردد.

ktable¶

فقط جریان های کافکا DSL مفهوم ktable را دارد.

KTable انتزاع یک جریان ChangeLog است ، که در آن هر رکورد داده نشانگر بروزرسانی است. به طور دقیق تر ، مقدار در یک ضبط داده به عنوان "به روزرسانی" آخرین مقدار برای همان کلید رکورد تعبیر می شود ، در صورت وجود (اگر یک کلید مربوطه هنوز وجود نداشته باشد ، به روزرسانی درج در نظر گرفته می شود). با استفاده از قیاس جدول ، یک ضبط داده در یک جریان ChangeLog به عنوان یک درج/بروزرسانی UPSERT ARKA تعبیر می شود زیرا هر ردیف موجود با همان کلید بازنویسی می شود. همچنین ، مقادیر تهی به روشی خاص تفسیر می شود: یک رکورد با مقدار تهی نشان دهنده "حذف" یا سنگ قبر برای کلید رکورد است.

برای نشان دادن ، تصور می کنیم دو سوابق داده زیر به جریان ارسال می شوند:

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

اثرات تراکم ورود به سیستم Kafka: روش دیگری برای تفکر در مورد Kstream و Ktable به شرح زیر است: اگر می خواهید یک Ktable را در یک موضوع Kafka ذخیره کنید ، احتمالاً می خواهید ویژگی تراکم ورود به سیستم Kafka را فعال کنید ، به عنوان مثال. برای صرفه جویی در فضای ذخیره سازی.

با این حال ، امکان ایجاد تراکم ورود به سیستم در مورد kstream ایمن نخواهد بود زیرا به محض اینکه تراکم ورود به سیستم شروع به پاکسازی سوابق داده های قدیمی تر از همان کلید می کند ، معنایی داده ها را می شکند. برای انتخاب دوباره مثال تصویر ، شما به طور ناگهانی 3 مورد را برای آلیس به جای 4 دریافت می کنید زیرا تراکم ورود به سیستم ("آلیس" ، 1) داده را حذف کرده است. از این رو ، تراکم ورود به سیستم برای یک KTable (جریان ChangeLog) کاملاً بی خطر است اما برای یک kstream (جریان ضبط) اشتباه است.

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

KTable همچنین توانایی جستجوی مقادیر فعلی سوابق داده را توسط کلیدها فراهم می کند. این قابلیت های ظاهر جدول از طریق عملیات پیوستن در دسترس است (همچنین به راهنمای توسعه دهنده) و همچنین از طریق نمایش داده های تعاملی مراجعه کنید.

globalktable¶

فقط جریان های کافکا DSL مفهوم GlobalKtable را دارد.

مانند KTable ، GlobalKtable انتزاعی از یک جریان ChangeLog است ، که در آن هر ضبط داده نشان دهنده یک به روزرسانی است.

یک GlobalKtable با داده هایی که در آن جمع می شوند با یک KTable متفاوت است ، یعنی کدام داده های مربوط به موضوع Kafka اساسی در جدول مربوطه خوانده می شود. کمی ساده ، تصور کنید که یک موضوع ورودی با 5 پارتیشن دارید. در برنامه خود می خواهید این موضوع را در یک جدول بخوانید. همچنین ، شما می خواهید برنامه خود را در 5 نمونه برنامه برای حداکثر موازی سازی اجرا کنید.

  • اگر موضوع ورودی را در ktable بخوانید ، نمونه "محلی" KTABLE از هر نمونه برنامه با داده هایی از تنها 1 پارتیشن 5 پارتیشن موضوع جمع می شود.
  • اگر موضوع ورودی را در یک GlobalKtable بخوانید ، نمونه محلی GlobalKtable از هر نمونه برنامه با داده های همه پارتیشن های موضوع جمع می شود.

GlobalKtable توانایی جستجو در مقادیر فعلی سوابق داده را توسط کلیدها فراهم می کند. این عملکرد جدول از طریق عملیات پیوستن (همانطور که در پیوستن به راهنمای توسعه دهنده توضیح داده شده است) و kafka جریان های تعاملی در دسترس است.

مزایای جداول جهانی:

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

پایین آمدن جداول جهانی:

  • افزایش مصرف محلی در مقایسه با (تقسیم بندی شده) KTable افزایش یافته است زیرا کل موضوع ردیابی می شود.
  • افزایش بار کارگزار شبکه و کافکا در مقایسه با (تقسیم بندی شده) KTable زیرا کل موضوع خوانده می شود.

یک جنبه مهم در پردازش جریان مفهوم زمان و نحوه مدل سازی و یکپارچه سازی آن است. به عنوان مثال ، برخی از عملیات مانند پنجره ها بر اساس مرزهای زمانی تعریف می شوند.

جریان کافکا از مفاهیم زیر از زمان پشتیبانی می کند:

زمان رویداد¶

نقطه زمانی که یک رویداد یا سابقه داده رخ داده است (یعنی در ابتدا توسط منبع ایجاد شده است). دستیابی به معناشناسی زمان رویداد به طور معمول نیاز به تعبیه زمانبندی در سوابق داده در زمان تولید رکورد داده دارد.

  • مثال: اگر این رویداد یک تغییر مکان جغرافیایی است که توسط یک سنسور GPS در یک ماشین گزارش شده است ، زمان رویداد مرتبط با آن زمان است که سنسور GPS تغییر مکان را ضبط کند.

زمان پردازش¶

نکته زمانی که اتفاق می افتد که رویداد یا سابقه داده توسط برنامه پردازش جریان پردازش می شود (یعنی وقتی سابقه مصرف می شود). زمان پردازش ممکن است میلی ثانیه ، ساعت ، روز و غیره باشد.

  • مثال: یک برنامه تحلیلی را تصور کنید که داده های مکان جغرافیایی گزارش شده از سنسورهای خودرو را برای ارائه آن به داشبورد مدیریت ناوگان می خواند و پردازش می کند. در اینجا ، پردازش زمان در برنامه تحلیلی ممکن است میلی ثانیه یا ثانیه باشد (مانند خطوط لوله در زمان واقعی بر اساس جریان های کافکا و کافکا) یا ساعت ها (مانند خطوط لوله دسته ای بر اساس Apache Hadoop یا Apache Spark) پس از زمان رویداد.

مصرف زمان

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

زمان جریان

حداکثر زمانی که تاکنون در تمام سوابق پردازش شده دیده می شود. Kafka Streams زمان جریان را بر اساس هر کار ردیابی می کند.

Timestamps¶

Kafka Streams از طریق به اصطلاح استخراج کننده های Timestamp یک جدول زمانی را به هر رکورد داده اختصاص می دهد. این جدول زمانی در هر رکورد پیشرفت یک جریان را با توجه به زمان توصیف می کند (اگرچه سوابق ممکن است خارج از سفارش در جریان باشد) و با عملیات وابسته به زمان مانند پیوستن به آنها اعمال می شود. ما آن را زمان رویداد برنامه می نامیم تا در هنگام اجرای این برنامه ، با دیوار ساعت متمایز شود. زمان رویداد همچنین برای همگام سازی چندین جریان ورودی در همان برنامه استفاده می شود.

پیاده سازی های بتن استخراج کننده های Timestamp ممکن است بر اساس محتوای واقعی سوابق داده مانند یک قسمت Timestamp تعبیه شده برای ارائه معناشناسی زمان رویداد یا زمان مصرف ، بازیابی یا محاسبه کنند ، یا از هر روش دیگری مانند بازگرداندن زمان کنونی دیوار کنونی استفاده کنند. زمان پردازش ، از این طریق معناشناسی زمان پردازش به برنامه های پردازش جریان می دهد. بنابراین توسعه دهندگان می توانند بسته به نیازهای تجاری خود ، مفاهیم/معناشناسی مختلفی را اجرا کنند.

سرانجام ، هر زمان که یک برنامه Kafka Streams سوابق را به کافکا می نویسد ، پس از آن به این سوابق جدید نیز اختصاص می یابد. نحوه اختصاص زمان بندی ها به متن بستگی دارد:

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

برای جمع آوری و پیوستن ، جدول زمانی با استفاده از قوانین زیر محاسبه می شود.

  • برای پیوستن (جریان جریان ، جدول جدول) که دارای سوابق ورودی چپ و راست است ، زمان بندی ضبط خروجی حداکثر (سمت چپ. tts ، راست. tts) اختصاص داده می شود.
  • برای پیوستن به جدول ، به رکورد خروجی زمان سنج از رکورد جریان اختصاص می یابد.
  • برای جمع آوری ، جریان های کافکا همچنین حداکثر زمان سنج را در تمام سوابق ، در هر کلید ، یا در سطح جهان (برای غیر پنجره) یا در هر پنجره محاسبه می کند.
  • به عملیات بدون تابعیت ، جدول زمانی رکورد ورودی اختصاص داده می شود. برای Flatmap و خواهران و برادرانی که چندین سوابق منتشر می کنند ، تمام سوابق خروجی زمان سنج را از سابقه ورودی مربوطه به ارث می برند.

جنبه های دیگر زمان

وقت خود را بدانید: هنگام کار با زمان ، باید اطمینان حاصل کنید که جنبه های اضافی زمان مانند مناطق زمانی و تقویم ها به درستی هماهنگ شده اند - یا حداقل درک و ردیابی شده - در کل خطوط لوله داده جریان شما. به عنوان مثال ، این کمک می کند تا در مورد مشخص کردن اطلاعات زمان در UTC یا در زمان UNIX (مانند ثانیه از زمان دوره) به توافق برسد. شما همچنین نباید موضوعات را با معناشناسی زمانی مختلف مخلوط کنید.

جمع بندی

یک عملیات تجمیع یک جریان یا جدول ورودی را به خود اختصاص می دهد و با ترکیب چندین سوابق ورودی در یک رکورد خروجی واحد ، جدول جدیدی را به دست می آورد. نمونه هایی از تجمع تعداد محاسبات یا مبلغ محاسبات است.

در جریان های Kafka DSL ، یک جریان ورودی از یک عملیات جمع آوری می تواند یک kstream یا ktable باشد ، اما جریان خروجی همیشه یک ktable خواهد بود. این به جریان های کافکا اجازه می دهد تا پس از تولید و انتشار مقدار ، مقدار کل را به روز کنند. هنگامی که چنین ورود خارج از سفارش اتفاق می افتد ، جمع آوری kstream یا ktable یک مقدار کل جدید را منتشر می کند. از آنجا که خروجی KTable است ، مقدار جدید در نظر گرفته می شود که مقدار قدیمی را با همان کلید در مراحل پردازش بعدی بازنویسی کند. برای کسب اطلاعات بیشتر در مورد سوابق خارج از سفارش ، به رسیدگی به خارج از سفارش مراجعه کنید.

پیوستن

یک عملیات پیوستن دو جریان ورودی و/یا جداول را بر اساس کلیدهای سوابق داده آنها ادغام می کند و یک جریان/جدول جدید را به دست می آورد.

عملیات پیوستن موجود در جریان های Kafka DSL بر اساس آن که نوع جریان ها و جداول در حال پیوستن است متفاوت است. به عنوان مثال ، Kstream-Kstream در مقابل پیوندهای Kstream-Ktable به هم می پیوندد.

پنجره

Windowing به شما امکان می دهد نحوه گروه بندی سوابق را که دارای همان کلید برای عملیات حالت مانند تجمع یا پیوستن به به اصطلاح ویندوز هستند ، کنترل کنید. ویندوز در هر کلید رکورد ردیابی می شود.

Windowing operations are available in the Kafka Streams DSL . When working with windows, you can specify a grace period for the window that indicates when window results are final. This grace period controls how long Kafka Streams will wait for out-of-order data records for a window. If a record arrives after the grace period of a window has passed (i.e., record.ts>Window-end-Time + Grace-Period) ، رکورد دور ریخته می شود و در آن پنجره پردازش نمی شود.

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

جریان های کافکا قادر به اجرای صحیح سوابق خارج از سفارش برای معناشناسی زمانی مربوطه (زمان رویداد) است.

دوره گریس در مقابل زمان نگهداری

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

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

نمایش داده شد

نمایش داده های تعاملی به شما امکان می دهد لایه پردازش جریان را به عنوان یک پایگاه داده تعبیه شده سبک وزن درمان کنید و مستقیماً از آخرین وضعیت برنامه پردازش جریان خود پرس و جو کنید. شما می توانید این کار را بدون نیاز به تحقق آن حالت در پایگاه داده های خارجی یا ذخیره خارجی انجام دهید.

نمایش داده های تعاملی معماری را ساده می کند و منجر به معماری بیشتر کاربردی محور می شود.

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

بهترین های هر دو جهان: البته شما همچنین می توانید معماری های ترکیبی را اجرا کنید که در آن به عنوان مثال ، برنامه شما ممکن است به صورت تعاملی پرسیده شود اما در عین حال برخی از نتایج خود را با سیستم های خارجی (به عنوان مثال از طریق Kafka Coect) به اشتراک می گذارد.

../_images/streams-interactive-queries-01.png

بدون نمایش داده های تعاملی: افزایش پیچیدگی و ردپای سنگین تر معماری.¶

../_images/streams-interactive-queries-02.png

با نمایش داده های تعاملی: معماری ساده تر و کاربردی بیشتر.¶

در اینجا برخی از نمونه های مورد استفاده برای برنامه هایی که از پرس و جوهای تعاملی بهره مند می شوند وجود دارد:

  • نظارت بر زمان واقعی: داشبورد جلویی که اطلاعات تهدید را فراهم می کند (به عنوان مثال ، سرورهای وب که در حال حاضر مورد حمله مجرمان سایبری قرار دارند) می توانند مستقیماً از یک برنامه جریان کافکا پرس و جو کنند که به طور مداوم با پردازش داده های شبکه از راه دور در زمان واقعی ، اطلاعات مربوطه را تولید می کند.
  • بازی ویدیویی: یک برنامه جریان کافکا به طور مداوم به روزرسانی های مکان را از بازیکنان در جهان بازی ردیابی می کند. سپس یک برنامه همراه موبایل می تواند مستقیماً از برنامه Kafka Streams پرس و جو کند تا موقعیت فعلی یک بازیکن را به دوستان و خانواده نشان دهد و از آنها دعوت کند تا همراه شوند. به طور مشابه ، فروشنده بازی می تواند از داده ها برای شناسایی نقاط مهم غیرمعمول بازیکنان استفاده کند ، که ممکن است نشانگر یک اشکال یا یک مسئله عملیاتی باشد.
  • ریسک و کلاهبرداری: یک برنامه کافکا به طور مداوم معاملات کاربر را برای ناهنجاری ها و رفتار مشکوک تجزیه و تحلیل می کند. یک برنامه بانکی آنلاین می تواند به طور مستقیم از برنامه Kafka Streams پرس و جو کند وقتی کاربر وارد سیستم می شود تا دسترسی به آن دسته از کاربرانی را که به عنوان مشکوک پرچم گذاری شده اند ، رد کند.
  • تشخیص روند: یک برنامه جریان Kafka به طور مداوم آخرین نمودارهای برتر را در ژانرهای موسیقی بر اساس رفتار گوش دادن به کاربر که در زمان واقعی جمع آوری می شود ، محاسبه می کند. برنامه های کاربردی موبایل یا دسک تاپ یک فروشگاه موسیقی می توانند در حالی که کاربران در حال مرور فروشگاه هستند ، از آخرین نمودارها به صورت تعاملی پرس و جو کنید.

برای اطلاعات بیشتر ، به راهنمای توسعه دهنده مراجعه کنید.

ضمانت های پردازش

Kafka Streams از ضمانت های پردازش در یک بارداری و دقیقاً یکپارچه پشتیبانی می کند.

سوابق معناشناسی در تنها یک باره هرگز از بین نمی روند اما ممکن است دوباره جمع شوند. اگر برنامه پردازش جریان شما ناکام باشد ، هیچ سوابق داده از بین نمی رود و پردازش نمی شود ، اما برخی از سوابق داده ممکن است دوباره خوانده شوند و بنابراین دوباره پردازش شوند. معنایی AT LEAST-ONCE به طور پیش فرض (پردازش. guarantee = "at_least_once") در پیکربندی جریان شما فعال می شود. سوابق معنایی دقیقاً یک بار یک بار پردازش می شوند. حتی اگر یک تولید کننده یک رکورد تکراری ارسال کند ، دقیقاً یک بار به کارگزار نوشته شده است. پردازش جریان دقیقاً یکپارچه ، امکان اجرای یک عملیات-فرآیند خواندن نوشتن دقیقاً یک بار است. تمام پردازش دقیقاً یک بار اتفاق می افتد ، از جمله پردازش و وضعیت مادی ایجاد شده توسط کار پردازش که به کافکا نوشته شده است. برای فعال کردن معناشناسی دقیق یکپارچه ، پردازش را تنظیم کنید.

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

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

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

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

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

رهگیرهای نظارتی همزمان را نمی توان در مرکز کنترل همخوانی در ارتباط با Exactly Once Semantics (EOS) پیکربندی کرد.

رسیدگی خارج از سفارش¶

علاوه بر تضمین اینکه هر رکورد دقیقاً یک بار پردازش می شود، یکی دیگر از مسائل چالش برانگیز که بسیاری از برنامه های پردازش جریانی با آن مواجه هستند، نحوه مدیریت داده های نامرتب است که ممکن است بر منطق تجاری آنها تأثیر بگذارد. در Kafka Streams، دو دلیل وجود دارد که به طور بالقوه می تواند منجر به دریافت نامتعادل داده ها با توجه به مهر زمانی آنها شود:

  • در یک پارتیشن موضوعی، مُهر زمانی یک رکورد ممکن است به طور یکنواخت همراه با افست های آنها افزایش پیدا نکند. از آنجایی که Kafka Streams همیشه سعی می کند رکوردها را به ترتیب افست پردازش کند، می تواند باعث شود رکوردهایی با مُهر زمانی بزرگ تر (اما افست های کوچک تر) زودتر از رکوردهایی با مهر زمانی کوچک تر (اما افست های بزرگ تر) در همان پارتیشن موضوع پردازش شوند.
  • یک کار جریانی ممکن است چندین پارتیشن موضوعی را پردازش کند، و اگر برنامه طوری پیکربندی شده باشد که منتظر نماند تا همه پارتیشن ها حاوی برخی داده های بافر باشند و از پارتیشنی با کمترین مهر زمانی انتخاب کند تا رکورد بعدی را پردازش کند، رکوردها بعداً برای پارتیشن های موضوع دیگر واکشی می شوند. ممکن است دارای مهرهای زمانی کوچک تر از رکوردهای پردازش شده باشد، که باعث می شود رکوردهای قدیمی تر پس از رکوردهای جدیدتر پردازش شوند. برای اطلاعات بیشتر، به max. task. idle. ms مراجعه کنید.

برای عملیات های بدون حالت، داده های نامرتب بر منطق پردازش تأثیر نمی گذارند، زیرا در هر زمان فقط یک رکورد در نظر گرفته می شود، بدون اینکه به تاریخچه رکوردهای پردازش شده گذشته نگاهی بیندازید.

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

اصطلاحات خارج از سفارش ¶

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

برای فعال کردن استفاده و درک مداوم از مفاهیم سفارش ، از تعاریف زیر استفاده کنید.

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

خواندن پیشنهادی

این وب سایت شامل محتوایی است که در بنیاد نرم افزار Apache تحت شرایط مجوز Apache V2 ایجاد شده است.

حساب اسلامي...
ما را در سایت حساب اسلامي دنبال می کنید

برچسب : نویسنده : کامران فیوضات بازدید : <-PostHit-> تاريخ : پنجشنبه 8 تير 1402 ساعت: 21:30