از JDK 24 به بعد میتوانید از API خاصی به نام Gatherers API برای مدلسازی عملیاتهای میانی در Stream API استفاده کنید. از نظر طراحی مشابه Collector API برای عملیاتهای پایانی Stream است.
چرا Stream API به چنین ویژگیای نیاز داشت؟ Stream API بسیار غنی است و راههای زیادی برای پردازش دادههای حافظهای بر اساس الگوی map-filter-reduce به شما میدهد. Stream API روی Spliterator API ساخته شده تا قابلیت موازیسازی عملیات را مدلسازی کند. تنوع این دو الگو امکانات زیادی میدهد. تنها نقطه ضعف این است که Spliterator API استفاده آسانی ندارد و به کد ساده و خوانا منجر نمیشود.
Gatherer API الگوهای سادهتری از Spliterator API با پشتیبانی عالی از موازیسازی ارائه میدهد. Gatherer API بر دو عنصر اصلی بنا شده: رابط Gatherer و کلاس factory Gatherers.
Gatherer شیئی است که میتوانید به متد gather() رابط Stream بدهید. متد gather() یک عملیات میانی Stream API است.
Gatherer API یک API ساده نیست و نباید از آن استفاده کنید مگر دلیل خوبی داشته باشید. حتی اگر بتواند سادهترین عملیات Stream مانند map، filter یا flat-map را مدلسازی کند، ساختن Gatherer برای این عملیات پیچیده خواهد بود.
Gatherer API برای ایجاد عملیات پیچیدهای که قبلاً در Stream API موجود نیستند طراحی شده. نمونهها:
آموختن نحوهی کار Gatherer API فرآیند ناامیدکنندهای است. باید از سری مثالهای سادهای عبور کنید.
Upstream: استریمی که عناصر از آن میآیند. Downstream: استریمی که Gatherer عناصر را به آن میفرستد.
نوشتن Gatherer دربارهی پیادهسازی رابط Gatherer است که رابط تابعی است. متد Gatherer.of() را با یک Integrator استفاده کنید.
میتوانید mapper ای بنویسید که نگاشت واقعی انجام میدهد:
result = [ONE, TWO, THREE]
مثالی که نوشتید integrator ای از نوع <Void, String, String>. نوع اول نوع state داخلی است که اینجا استفاده نمیکنیم. نوع دوم نوع عناصر مصرفشده از upstream است و نوع سوم نوع عناصر تولیدشده به downstream است.
توجه کنید این متد یک boolean برمیگرداند. خوشبختانه متد downstream.push() نیز boolean برمیگرداند، به همین دلیل میتوانید مثال را کامپایل و اجرا کنید. این boolean بسیار مهم است و نقش دقیق آن بعداً در این بخش پوشش داده میشود.
اکنون میتوانید mapper ای بنویسید که نگاشت واقعی انجام میدهد (به جای تابع همانی). بیایید این gatherer را در یک متد قرار دهیم تا بتوانید mapper خود را به عنوان آرگومان بدهید.
result = [one, two]
اجرای این کد نتیجه زیر را تولید میکند.
result = [ONE, TWO, THREE]
میتوانید تابع در خط ۸ را به مورد زیر تغییر دهید. توجه کنید به لطف استفاده از var نیازی به تغییر نوع gatherer تولیدشده نیست.
result = [o, n, e, t, w, o, t, h, r, e, e]
نتیجه سپس به صورت زیر در میآید.
result = [3, 3, 5]
سه باگ در flatMapper.apply(mapped).forEach(downstream::push) وجود دارد:
Downstream thread-safe نیست. باید از flatMapper.apply(mapped).sequential() استفاده کنید.چندین عملیات میانی از state قابل تغییر استفاده میکنند:
limit(): شمارش عناصر پردازششدهdistinct(): ذخیره عناصر دیدهشده در HashSetdropWhile(): تغییر state داخلیبرای Gatherer با state، به یک initializer (Supplier) نیاز دارید و باید از Gatherer.ofSequential() استفاده کنید.
result = [4, 5, 4, 3, 2, 1]
اجرای این کد نتیجه زیر را تولید میکند.
result = [o, n, e, t, w, o, t, h, r, e, e]
توجه کنید خط ۵ این کد: return true بسیار مشکوک به نظر میرسد. و در واقع، این روش مدیریت boolean بازگرداندهشده توسط فراخوانی downstream::push نیست. بعداً در این بخش خواهید دید چرا این کد اشتباه است و چگونه آن را اصلاح کنید.
همچنین توجه کنید حداقل دو باگ ظریف در نوشتن flatMapper.apply(element).forEach(downstream::push) وجود دارد که بعداً بحث خواهیم کرد. با این حال، استفاده از این کد در مثال قبلی مشکلی ندارد.
سه اصل با state rejecting downstream وجود دارد:
اجرای این کد نتیجه زیر را تولید میکند.
result = [O, N, E, T, W, O]
یک بار دیگر، این مثال را با احتیاط بررسی کنید، چون return true در خط ۱۲ کد صحیح نیست. بیشتر درباره آن بعداً در این بخش خواهیم گفت.
چرا این کد سریعتر اجرا میشود؟ چون اشیاء کمتری نسبت به کد استریم معادل ایجاد میکند. هر فراخوانی متد استریم یک شیء Stream جدید ایجاد میکند که overhead محسوب میشود. در این کد فقط دو شیء Stream ایجاد شده که میتواند یک بهینهسازی باشد.
بعداً در این بخش دربارهی زنجیرهای کردن Gatherer ها بیشتر خواهید آموخت که راه بهتری (و توصیهشده!) برای ادغام چندین عملیات Gatherer در یکی ارائه میدهد. در هر صورت، قبل از استفاده از آن باید بسنجید آیا این نوع کد در محیط production سودی برایتان دارد.
توجه کنید باگ ظریفی که در مثال قبلی ذکر کردیم هنوز در نوشتن flatMapper.apply(mapped).forEach(downstream::push) وجود دارد. پس حالا زمان اصلاح آن است.
حداقل دو باگ در کد زیر وجود دارد: flatMapper.apply(mapped).forEach(downstream::push)، و یک باگ سوم هم هست.
باگ اول این است که نمونهی Downstream یک شیء thread-safe نیست. فراخوانی forEach(downstream::push) یکی از قوانین Stream API را نقض میکند: lambda هایی که در Stream استفاده میکنید نباید effect های جانبی داشته باشند. اگر از این قانون پیروی نکنید، باید با استریمهای موازی بسیار محتاط باشید: میتوانند race condition ایجاد کنند، که دقیقاً همین مورد اینجاست. راهحل این است که از flatMapperای که ممکن است استریم موازی برگرداند و race condition در Downstream ایجاد کند، محافظت کنید. میتوانید این کار را با فراخوانی sequential روی استریم بازگرداندهشده انجام دهید: flatMapper.apply(mapped).sequential(). اگر استریم از قبل sequential باشد، این فراخوانی کاری انجام نمیدهد و this برمیگرداند، و اگر استریم موازی باشد، آن را غیرموازی میکند.
Files.lines()) باز شده که فراخوانی نکردن آن منجر به نشت منبع میشود. چون Stream رابط AutoCloseable را پیادهسازی میکند، میتوانید آن را در یک عبارت try-with-resources باز کنید. الگو به صورت زیر در میآید.
result = [1, 2, 3]
و یکی دیگر هم هست که احتمالاً باید از آن محافظت کنید. دوست قدیمی ما NullPointerException معروف که هیچکس انتظارش را ندارد. ممکن است کسی flat-mapperای به متد شما بدهد که null برگرداند. مقابله با چنین مشکلی به برنامهی شما بستگی دارد. شاید باگی باشد که باید گزارش دهید و در آن صورت میتوانید تصمیم بگیرید سریعاً fail کنید و این NullPointerException را دوباره throw کنید. یا میتوانید آن را نادیده بگیرید و از آنجا که چیزی برای ارسال به Downstream وجود ندارد، کاری انجام ندهید.
این بخش آخر مربوط به نحوهی مدیریت boolean بازگرداندهشده توسط downstream::push است، پس الگوی کامل و صحیح را بعداً در این بخش خواهید دید.
دومین کاری که Gatherer میتواند انجام دهد مدیریت state داخلی قابل تغییر است. چندین عملیات میانی برای کار در Stream API از state قابل تغییر استفاده میکنند.
limit() باید تعداد عناصری که قبل از قطع استریم پردازش کرده را بشمارد.distinct() باید عناصری که قبلاً دیده را در HashSet ذخیره کند.dropWhile() باید state داخلی را تغییر دهد تا بداند باز است یا نه.اصلی در Stream API وجود دارد که میگوید هرگز نباید state خارجی قابل تغییری را از داخل استریم تغییر دهید. دو دلیل برای این موضوع وجود دارد. اول: این کار ممکن است مانع اعمال برخی بهینهسازیها روی استریمها شود. دوم: از استفاده از استریمهای موازی جلوگیری میکند، چون اگر محتاط نباشید race condition ایجاد میشود.
بیایید ابتدا یک رفتار شبیه dropWhile() با یک gatherer ایجاد کنیم.
شاید متوجه شده باشید که پارامتر اول integrator هایی که در بخش قبلی نوشتیم را نادیده گرفتیم. خب، حالا زمان در نظر گرفتن آنهاست. Gatherer ای که state داخلی قابل تغییر را مدیریت میکند باید این state فعلی را بداند تا تصمیم بگیرد چه کند. پس integrator این state را به عنوان پارامتر اول میگیرد. چون این state قابل تغییر است، integrator میتواند آن را تغییر دهد و اگر این کار را انجام دهد، دفعهی بعدی که فراخوانی شود این تغییر را میبیند.
اما این همهی ماجرا نیست: به راهی برای ایجاد نمونهی اولیهی این state قابل تغییر نیز نیاز دارید.
اگر به یک collection قابل تغییر نیاز دارید، میتوانید از پیادهسازی ارائهشده توسط JDK استفاده کنید. لازم نیست thread-safe باشد، پس ArrayList کافی است. اگر به یک شمارنده یا boolean نیاز دارید، باید wrapper قابل تغییر خودتان را روی این شمارنده ایجاد کنید. استفاده از متغیر atomic زیادهروی است، چون به thread safety نیاز ندارید. در این مورد، داشتن فقط یک integrator کافی نیست: به supplier ای نیاز دارید که پیادهسازی Gatherer بتواند آن را برای ایجاد این container قابل تغییر اولیه فراخوانی کند. در Gatherer API، این supplier initializer نامیده میشود.
بیایید gatherer ای ایجاد کنیم که میتواند دروازه را باز کند تا عناصر جریان یابند هنگامی که predicate نادرست شود.
result = [1, 2, 3, 4, 5]
اگر به downstream بگویید که اهمیت نمیدهید و به ارسال عناصر ادامه دهید، در واقع اتفاقی نمیافتد. عناصری که ارسال میکنید به آرامی نادیده گرفته میشوند. Exception ای ایجاد نمیشود.
دو متد factory روی رابط Integrator وجود دارد.
تعریف Integrator به عنوان greedy به این معناست که هرگز به تنهایی stream را قطع نمیکند. Gatherer API میتواند برخی بهینهسازیها را فعال کند.
اجرای کد قبلی خروجی زیر را چاپ میکند.
result = [4, 5, 4, 3, 2, 1]
اولین عنصری که در این متد dropWhile() میبینید، تعریف box قابل تغییر شماست. این یک کلاس محلی ساده است که میتوانید آن را در متد خود تعریف کنید. در مرحلهی بعدی خواهید دید چگونه میتوانید آن را به کلاس ناشناس تبدیل کنید. سپس میتوانید یک Supplier بنویسید تا نمونهای از کلاس Box ایجاد کند. این supplier فقط یک بار توسط پیادهسازی gatherer شما فراخوانی میشود.
پیادهسازی integrator شما متفاوت است و نمونهی Box ایجادشده را دریافت میکند. چون قابل تغییر است، هر کاری که لازم است با آن انجام دهید. اینجا، فقط boolean پیچیدهشده در آن را هنگام true شدن predicate تغییر میدهیم و در آن صورت عناصر را به downstream ارسال میکنیم. توجه کنید این integrator همیشه true برمیگرداند، که در واقع کار صحیحی نیست. بعداً آن را اصلاح خواهیم کرد.
بزرگترین تفاوت این است که اکنون باید به جای متد factory ساده Gatherer.of() از متد factory Gatherer.ofSequential() استفاده کنید. استفاده از ofSequential() به جای of() به API میگوید که این gatherer نمیتواند به صورت موازی استفاده شود. بعداً در این بخش دربارهی gatherer های موازی و اینکه چرا نمیتوان این gatherer را به صورت موازی استفاده کرد بیشتر خواهید آموخت.
میتوانید این gatherer را به روش دیگری با استفاده از کلاسهای ناشناس و بهرهبرداری از انواع non-denotable بنویسید.
result = [ONE, TWO, THREE, DONE]
اجرای کد قبلی همان نتیجهی قبل را به شما میدهد.
result = [4, 5, 4, 3, 2, 1]
این کد کار میکند چون کلاس ناشناس به عنوان نوع non-denotable کامپایل میشود، که تا زمانی که این متغیر را در متغیری با نوع غیرقابل استنباط قرار ندهید حفظ میشود. اینجا، نوع آرگومان box در خط ۴ توسط کامپایلر استنباط میشود، این نوع حفظ میشود و میتوانید به فیلد کلاس ناشناس خود دسترسی داشته باشید.
با این حال، اگر نوع صریحی ارائه دهید، نوع non-denotable توسط کامپایلر استفاده نمیشود و از بین میرود. استفاده از کلاسهای محلی متد ممکن است خواناتر باشد و همیشه کار میکند.
بیایید مثال دومی را بررسی کنیم که شامل پیادهسازی رفتاری شبیه distinct() است.
هر بار که عنصری مصرف میکنید، باید بررسی کنید آیا قبلاً دیده شده یا نه. و در آن صورت، نباید آن را به downstream ارسال کنید. میتوانید از یک HashSet معمولی به عنوان state قابل تغییر خود استفاده کنید، بدون نیاز به بستهبندی یا thread-safe کردن آن، چون در یک thread واحد استفاده خواهد شد.
یکی از ویژگیهای خوب متد Set.add() این است که اگر شیء اضافه شده باشد true برمیگرداند، یعنی قبلاً در مجموعه نبوده، و اگر اضافه نشده باشد false برمیگرداند. پس میتوانید در این مورد از این ویژگی استفاده کنید.
بیایید این Gatherer را ایجاد کنیم.
result = [one, six, two, five, four, seven, three]
اجرای مثال قبلی خروجی زیر را چاپ میکند.
result = [1, 2, 3, 4, 5]
یک بار دیگر، return true در خط ۸ کد صحیح نیست، بعداً آن را اصلاح خواهید کرد.
قطع استریم کاری است که توسط چندین عملیات میانی انجام میشود، از جمله limit() و takeWhile().
فرض کنید لیستی با یک میلیون عنصر دارید و کد زیر را فراخوانی میکنید.
result = [O, N, E, T, W, O]
Gatherer API همیشه finisher را پس از اتمام ارسال عناصر به integrator فراخوانی میکند.
Gatherer از Gatherer.of() میتواند موازی باشد. Gatherer از Gatherer.ofSequential() فقط ترتیبی است.
حتی Gatherer ترتیبی در استریم موازی، از عملکرد موازی بخشهای قبل و بعد بهره میبرد. این ویژگی منحصربهفرد Gatherer API است.
سومین الگو Gatherer.of(initializer, integrator, combiner, finisher) state را بین threadها اشتراک نمیگذارد و combiner برای ادغام state ها نیاز دارد.
result = [one, six, two, five, four, seven, three]
این کد به درستی کار میکند، اما بهترین چیزی نیست که میتوانید بنویسید. سه اصل در مورد state rejecting downstream وجود دارد.
با رعایت این سه اصل، خطوط اول integrator شما بیفایده میشوند.
result = [THREE, FOUR, FIVE, SEVEN]
آخرین عنصر Stream API که باید بشناسید کلاس Gatherers است. پنج متود factory دارد.
fold() عملیاتی مشابه reduce() اما بدون محدودیتهای آن اجرا میکند.
result = {1234}
اجرای این کد نتیجه زیر را تولید میکند.
result = [1, 2, 3]
در خط ۱۰، یک عنصر دیگر ارسال میکنید، مقدار بازگرداندهشده توسط فراخوانی downstream.push() را دریافت میکنید و فوراً آن را به upstream خود منتقل میکنید. اگر boolean بازگرداندهشده false باشد، این integrator دیگر فراخوانی نمیشود.
سپس، اگر به خط 12 برسید، به این معناست که عناصر کافی به downstream ارسال کردهاید و کار تمام شده. پس در هر صورت false برمیگردانید.
با در نظر گرفتن این موضوع، اکنون میتوانید مثالهای قبلی که با return true به عنوان placeholder رها شده بودند را اصلاح کنید. نسخهی صحیح مثال Distinct به صورت زیر است، با فرض اینکه downstream شما نمیتواند در state ردکننده شروع شود.
result = [{1}, {12}, {123}, {1234}]
mapConcurrent() عناصر را با virtual thread ها به صورت همزمان نگاشت میکند. تعداد thread های فعال با maxConcurrency کنترل میشود.
دو استراتژی پنجره وجود دارد:
اجرای این مثال نتیجه زیر را تولید میکند.
result = [1, 2, 3, 4, 5]
از اصول Gatherer Limit پیروی میکند.
downstream.push() را منتقل کنید. اگر false برگردانید، integrator شما دیگر فراخوانی نمیشود.true برگردانید. نیازی به بررسی state ردکنندهی downstream نیست، چون میدانید false است. در واقع، integrator شما در حال حاضر فراخوانی شده، پس باید false باشد و downstream.push() را فراخوانی نکردهاید، پس نمیتواند تغییر کند.در هر صورت، باید به یاد داشته باشید که وقتی false برمیگردانید، نباید تصمیم بگیرید دوباره true برگردانید. به عنوان downstream استریم upstream، باید از قانونی پیروی کنید که downstream هرگز نمیتواند از state غیرپذیرنده به state پذیرنده تغییر کند.
به یاد داشته باشید وقتی false برمیگردانید، integrator شما دیگر فراخوانی نمیشود.
اگر نادیده بگیرید downstream چه میگوید و به ارسال عناصر ادامه دهید چه اتفاقی میافتد؟ در واقع هیچ چیز. عناصری که ارسال میکنید به آرامی نادیده گرفته میشوند. Exception ای ایجاد نمیشود. فقط هدر دادن منابع و overhead ای است که میتوانستید از آن اجتناب کنید. در برخی موارد این overhead میتواند بسیار پرهزینه باشد.
دو متد factory روی رابط Integrator وجود دارد.
Integrator.of() که یک Gatherer.Integrator ساده به عنوان پارامتر میگیرد.Integrator.ofGreedy() که یک Gatherer.Integrator.Greedy به عنوان پارامتر میگیرد.دو چیز وجود دارد که باید به یاد داشته باشید:
Gatherer.Integrator.Greedy یک extension از Gatherer.Integrator است.پس میتوانید موارد زیر را بنویسید. این دو integrator کاملاً معادل هم هستند.
result = [[one, two], [three, four], [five]]
result = [[1, 2, 3], [2, 3, 4], [3, 4, 5], [4, 5, 6], [5, 6, 7], [6, 7, 8], [7, 8, 9]]
میتوانید Gatherer های پنجرهای را با Gatherer های دیگر زنجیرهای کنید:
علاوه بر این، این چهار integrator به یک شکل کار میکنند. از نظر رفتاری، تفاوتی بین آنها وجود ندارد.
تعریف integrator به عنوان greedy به این معناست که این integrator هرگز به تنهایی stream را قطع نمیکند. همیشه state ردکنندهی downstreamای که عناصر را به آن ارسال میکند را منتقل میکند، یا اگر چیزی ارسال نمیکند true. پس Gatherer API میتواند این موضوع را در نظر بگیرد و برخی بهینهسازیها را برای اجرای gatherer شما فعال کند. اگر از قرارداد integrator greedy پیروی نکنید، کد شما ممکن است خطا دهد. اگر integrator سادهای با رفتار greedy دارید، کد شما خطا نمیدهد اما ممکن است برخی فرصتهای بهینهسازی را از دست بدهید.
پس در مرحلهی اول میتوانید این رفتار greedy را نادیده بگیرید و به integrator های ساده پایبند بمانید. پس از راهاندازی pipeline پردازش دادهها، میتوانید هر مرحله را بازبینی کنید و integrator های مربوطه را greedy کنید.
مواردی وجود دارد که gatherer شما باید منتظر تمام عناصر upstream خود بماند تا ارسال عناصر به downstream را شروع کند. این مورد وقتی صدق میکند که نیاز به ایجاد gathererای برای مرتبسازی عناصر upstream خود دارید. نمیتوانید تولید عناصر را شروع کنید مگر همهی آنها را دیده باشید.
چنین gatherer ای باید تمام عناصری که میبیند را به یک بافر داخلی اضافه کند. سپس وقتی عناصر دیگری به integrator ارسال نشد، این عناصر را به downstream ارسال کند. تا اینجا هیچ راهحلی برای این کار ندیدهاید. هیچ چیزی به gatherer شما نمیگوید که عناصر دیگری از upstream نمیآیند. در واقع، upstream عناصر را به integrator gatherer شما ارسال میکند و وقتی عناصر دیگری نیست، اتفاقی نمیافتد.
به همین دلیل Gatherer API عنصر دیگری به شما میدهد، به نام finisher، که توسط پیادهسازی هنگامی که عناصر دیگری به integrator شما ارسال نمیشود فراخوانی میشود.
این finisher را میتوان یک integrator سادهشده دید.
پس finisher دو پارامتر میگیرد: state شما و downstream، و چیزی برنمیگرداند. توسط BiConsumer<State, Downstream> مدلسازی شده است.
بیایید finisher سادهای بنویسیم، در موردی که gatherer شما هیچ state داخلی قابل تغییری تعریف نمیکند.
اینجا یک mapping gatherer تعریف میکنیم و به دلیلی نیاز داریم "DONE" در انتهای استریم داشته باشیم. finisher شما یک biconsumer است، اما چون این gatherer هیچ stateای تعریف نمیکند، پارامتر اول از نوع Void است و نادیده گرفته میشود.
result = [2.0, 3.0, 4.0, 5.0, 6.0, 7.0, 8.0]
اجرای کد قبلی نتیجه زیر را تولید میکند.
result = [ONE, TWO, THREE, DONE]
مرتبسازی یک استریم نیاز به اضافه کردن تمام عناصر به بافر قبل از شروع ارسال آنها به downstream دارد. این بافر همان state داخلی قابل تغییر gatherer شماست. نقش integrator ذخیرهی تمام عناصر در این بافر است. ارسال عناصر سپس توسط finisher انجام میشود.
سپس میتوانید کد زیر را بنویسید. توجه کنید از TreeSet استفاده میکنیم تا کد سادهتر شود، که اثر اضافهی حذف تکرارها از استریم شما را نیز دارد.
اجرای این کد خروجی زیر را چاپ میکند.
result = [one, six, two, five, four, seven, three]
چند نکته در این gatherer ارزش توجه دارد.
TreeSet اضافه میشوند تا مرتب باقی بمانند. همچنین تکرارها را حذف میکند.downstream.push() هرگز فراخوانی نمیشود، این downstream در state اصلی خود یعنی non-rejecting باقی میماند.downstream.push() بررسی کند آیا downstream هنوز عناصر را میپذیرد. استفاده از takeWhile() اینجا بسیار مفید است.توجه کنید از الگوی زیر برای ارسال عناصر به downstream استفاده میکنیم. در واقع، کاری که باید انجام دهید این است: در حالی که در state non-rejecting است، عناصر را به downstream ارسال کنید. این را به لطف بیانگری Stream API میتوانید به راحتی به این کد ترجمه کنید.
همچنین میتوانید از الگوی زیر استفاده کنید که همان کار را به روش کارآمدتری انجام میدهد.
همچنین میتوانید از مثال قبلی برای اصلاح نهایی flat-mapper که قبلاً دربارهاش صحبت کردیم استفاده کنید. کد کامل و صحیح که با استریمهای موازی، بستن stream و NullPointerException احتمالی سروکار دارد به صورت زیر است.
اجرای این کد نتیجهی مشابهی تولید میکند.
result = [O, N, E, T, W, O]
Gatherer API همیشه finisher شما را پس از اتمام ارسال عناصر به integrator فراخوانی میکند. ممکن است مواردی وجود داشته باشد که downstream به state ردکننده تغییر کرده باشد در حالی که integrator شما عناصر را به آن ارسال میکرده است. پس ممکن است انتخاب کنید ابتدا state ردکنندهی downstream خود را بررسی کنید، تا مطمئن شوید محاسبات پرهزینهای در finisher خود بدون دلیل انجام نمیدهید. بررسی state ردکنندهی downstream در integrator شما بیفایده است، اما ممکن است در finisher مفید باشد.
یکی از ویژگیهای شگفتانگیز Stream API این است که میتوانید محاسبات خود را در تمام هستههای CPU خود توزیع کنید فقط با فراخوانی parallel() روی یک استریم موجود.
Gatherer ها از استریمهای موازی به دو صورت پشتیبانی میکنند.
هنگامی که Gatherer ای با یکی از متدهای factory Gatherer.of() ایجاد میکنید، gathererای میسازید که میتواند به صورت موازی فراخوانی شود.
از طرف دیگر، هنگامی که از Gatherer.ofSequential() استفاده میکنید، gathererای میسازید که از موازیسازی پشتیبانی نمیکند.
به چه معناست که به صورت موازی فراخوانی شود، یا از موازیسازی پشتیبانی نکند؟ در Stream API، اگر تصمیم بگیرید استریم خود را موازی کنید، تمام عملیاتهای استریم شما به صورت موازی در threadهای مختلف اجرا میشوند. به زبان ساده: یک استریم موازی منبع دادههای شما را به چند بخش تقسیم میکند و هر بخش در thread خودش، روی هستههای مختلف CPU شما پردازش میشود. پس از پردازش همهچیز، نتایج جزئی دارید که باید با هم ادغام شوند. پس هر عملیات میانی به صورت موازی اجرا میشود، بدون نیاز به هماهنگسازی با آنچه در threadهای دیگر اتفاق میافتد. استثناهایی وجود دارد، چون برخی عملیاتهای میانی نیاز به هماهنگسازی دارند. این مورد دربارهی limit()، takeWhile()، distinct()، sorted() و برخی دیگر صدق میکند.
این موضوع در مورد Gatherer API صدق نمیکند. اگر از gatherer در یک استریم موازی استفاده کنید، دو اتفاق میافتد.
Gatherer.ofSequential() ایجاد کردهاید و چنین gatherer هایی فقط به صورت ترتیبی قابل استفاده هستند. در آن صورت، تمام عملیاتها قبل از فراخوانی gather() شما به صورت موازی اجرا میشوند، درست مانند یک استریم موازی عادی، سپس gatherer شما در یک thread اجرا میشود و عناصر را به یک downstream ارسال میکند و بقیهی استریم دوباره به صورت موازی اجرا میشود.پس حتی اگر Gatherer ای دارید که به هر دلیلی از موازیسازی پشتیبانی نمیکند، همچنان میتوانید از عملکرد ارائهشده توسط استریمهای موازی بهرهمند شوید. این ویژگی منحصربهفرد Gatherer API است.
سه الگو برای انتخاب جهت ایجاد Gatherer موازی دارید.
دو الگوی اول معادل الگوهای Gatherer های ترتیبی هستند.
Gatherer.of(integrator). این الگو را قبلاً در این بخش دیدهاید. Gatherer شما به هیچ state داخلی قابل تغییری اتکا نمیکند و به هیچ finisherای نیاز ندارد. نسخهی ترتیبی این متد وجود دارد: Gatherer.ofSequential(integrator).Gatherer.of(integrator, finisher). این gatherer به state قابل تغییری نیاز ندارد و یک finisher دارد. این هم الگویی است که قبلاً در این بخش دیدهاید. نسخهی ترتیبی این متد نیز وجود دارد: Gatherer.ofSequential(integrator, finisher).Gatherer.of(initializer,integrator,combiner,finisher). این الگویی است که اکنون پوشش خواهیم داد.الگوی سوم state قابل تغییری را اعلام میکند. در Gatherer API این state قابل تغییر بین threadهای مختلف استریم موازی شما اشتراکگذاری نمیشود. پس هر بار که Stream API gatherer را در یک thread اجرا میکند، نمونهی جدیدی از این state قابل تغییر برای کار با این gatherer ایجاد میکند. از یک طرف کار شما را آسانتر میکند، چون نیازی به نگرانی دربارهی state قابل تغییر thread-safe ندارید، چون هرگز اشتراکگذاری نمیشود.
اما از طرف دیگر، به راهی برای ادغام این نمونههای مختلف با هم نیاز دارید، به یک نمونهی واحد، پس از اتمام محاسبات. این موضوع متد factory مربوطه را کمی پیچیدهتر میکند، چون باید پارامتر چهارمی به نام combiner ارائه دهید که بتواند دو نمونه از state قابل تغییر شما را ترکیب کند.
پس اکنون چهار پارامتر برای این متد factory وجود دارد.
Supplier است.BinaryOperator مدلسازی شده. میتواند یکی از دو نمونه را برگرداند یا نمونهی جدیدی ایجاد کند. این combiner بسته به تعداد دفعاتی که منبع دادههای شما تقسیم شده، هر تعداد بار میتواند فراخوانی شود.BiConsumerای است که قبلاً دربارهاش صحبت کردیم. این finisher فقط یک بار، در انتهای محاسبات شما فراخوانی میشود.اگر به دلیلی به state داخلی قابل تغییر نیاز دارید اما نمیتوانید combiner ایجاد کنید، Stream API نمیتواند gatherer شما را به صورت موازی اجرا کند، چون نمیتواند نمونههای این state قابل تغییر را با هم ترکیب کند. در آن صورت باید از یک sequential gatherer استفاده کنید.
بیایید Gatherer مرتبساز خود را به Gatherer موازی تبدیل کنیم.
اجرای کد قبلی خروجی زیر را چاپ میکند.
result = [one, six, two, five, four, seven, three]
Gatherer API متد andThen() را ارائه میدهد. میتوانید gatherer های موجود را با این متد زنجیرهای کنید تا gatherer های بیشتری ایجاد کنید.
میتوانید زنجیرهای کردن را در برخی Gatherer هایی که قبلاً در این بخش نوشتید ببینید.
اجرای مثال قبلی خروجی زیر را چاپ میکند.
result = [THREE, FOUR, FIVE, SEVEN]
همانطور که تصور میکنید، محدودیتهایی روی Gatherer هایی که میخواهید زنجیرهای کنید وجود دارد. نوع عناصر تولیدشده توسط اولی باید با نوع عناصر پذیرفتهشده توسط دومی مطابقت داشته باشد.
آخرین عنصر Stream API که باید بشناسید کلاس Gatherers است. پنج متود factory دارد که gatherer های آمادهی استفادهای ایجاد میکنند که میتوانید در برنامهی خود استفاده کنید.
متد fold() یک عملیات reduce را از چپ به راست پیادهسازی میکند، مفید برای مواردی که نمیتوانید combiner بنویسید.
Folding شبیه کاهش یک استریم به نظر میرسد. اما عملیات کاهش نیاز دارد عملگر شما چندین ویژگی داشته باشد، از جمله associativity و در برخی موارد، عنصر همانی. وقتی عملیات کاهش شما عنصر همانی ندارد، متد معادل reduce() optional برمیگرداند که در صورت کاهش یک استریم خالی خالی است.
Folding این محدودیتها را ندارد. Folding فقط یک عملگر (در این مورد BiFunction) روی عناصر شما اعمال میکند تا نتیجه تولید کند. همچنین یک Supplier برای شروع فرآیند میگیرد.
توجه کنید این gatherer folding شما را به عنوان یک عملیات میانی پیادهسازی میکند. پس استریمی دارید که فقط میتواند یک عنصر تولید کند. اگر نیاز دارید این عنصر را بگیرید، باید findFirst().orElseThrow() را روی آن فراخوانی کنید. اگر نیاز دارید دادهها را بیشتر پردازش کنید، شاید استفاده از collector راهحل خوبی باشد.
میتوانید کد زیر را اجرا کنید تا این Gatherer را در عمل ببینید.
اجرای کد قبلی خروجی زیر را چاپ میکند.
result = {1234}
متد factory scan() gathererای ایجاد میکند که با دو عنصر کار میکند.
BiFunction.هر بار که upstream عنصر جدیدی به این gatherer ارسال میکند، BiFunction را روی عنصر قبلی و این عنصر بعدی از upstream اعمال میکند و سپس نتیجه را به downstream ارسال میکند.
میتوانید مثالی مشابه مثال fold بسازید. حتی اگر دو gatherer fold و scan مشابه به نظر برسند، یک تفاوت اصلی وجود دارد: gatherer fold فقط آخرین عنصر را به downstream ارسال میکند. پس به عنوان reducer در یک عملیات میانی عمل میکند. Gatherer scan تمام عناصر را هنگام پردازش ارسال میکند.
اجرای کد قبلی خروجی زیر را چاپ میکند.
result = [{1}, {12}, {123}, {1234}]
متد factory mapConcurrent() gathererای ایجاد میکند که میتواند دادههای شما را به صورت همزمان نگاشت کند. gatherer تولیدشده موازی نیست، چون به صورت داخلی با متد factory Gatherer.ofSequential() ایجاد شده.
این gatherer تمام عناصری که میتواند را مصرف میکند و سپس آنها را با Functionای که به عنوان آرگومان میدهید نگاشت میکند. هر نگاشت در یک virtual thread اجرا میشود که به طور خاص برای نگاشت این عنصر راهاندازی شده. تعداد virtual threadهای فعال در هر لحظه توسط آرگومان maxConcurrency کنترل میشود.
دو متد factory آخر Gatherer هایی ایجاد میکنند که با پنجرهها کار میکنند. پنجره یک بازه از اندیسها روی استریم شماست. فراخوانی یک windowing gatherer روی استریمی که مرتب نیست (در مفهوم ORDERED) بیمعنی است، چون چندین اجرا روی دادههای یکسان میتواند نتایج متفاوت و ناسازگار ایجاد کند.
راههای زیادی برای ایجاد چنین پنجرههایی وجود دارد. API دو مورد به شما میدهد.
[0,1,2], [3,4,5], [6,7,8], .... پس هیچ عنصری از یک پنجره به پنجرهی بعدی تکرار نمیشود. آخرین پنجره میتواند ناقص باشد، چون تعداد عناصری که upstream میتواند تولید کند ممکن است دقیقاً بر اندازهی پنجرهی مورد نیاز شما بخشپذیر نباشد.[0,1,2], [1,2,3], [2,3,4], .... تمام پنجرههای این استریم طول یکسانی دارند.متد factory windowFixed(windowSize) gathererای ایجاد میکند که عناصر ارسالشده توسط upstream را در لیستهای غیرقابل تغییر با اندازهی windowSize ذخیره میکند. اندیس اولین عنصر هر لیست با windowSize افزایش مییابد. پس یک عنصر مشخص از upstream فقط یک بار در یک لیست غیرقابل تغییر ذخیره میشود. هر لیست وقتی آماده شد به downstream این gatherer ارسال میشود. لیست آخر تمام عناصر باقیماندهی ارسالشده توسط upstream را شامل میشود. اگر خوششانس باشید تعداد درست (windowSize) خواهد بود، اما کد شما نباید به آن اتکا کند.
بیایید مثال زیر را بنویسیم تا این Gatherer را در عمل ببینیم.
اجرای کد قبلی خروجی زیر را چاپ میکند.همانطور که میبینید آخرین پنجره کوتاهتر از بقیه است، چون این استریم فقط پنج عنصر تولید میکند.
result = [[one, two], [three, four], [five]]
متد factory windowSliding(windowSize) gathererای ایجاد میکند که عناصر ارسالشده توسط upstream را در لیستهای غیرقابل تغییر با اندازهی windowSize ذخیره میکند. اندیس اولین عنصر هر لیست با ۱ افزایش مییابد. پس یک عنصر مشخص از upstream در windowSize لیست غیرقابل تغییر ذخیره میشود. این قانون برای آخرین عناصر upstream صدق نمیکند. هر لیست وقتی آماده شد به downstream این gatherer ارسال میشود. لیست آخر تعداد عناصر یکسانی با سایر لیستها دارد. پس آخرین عنصر ارسالشده توسط upstream فقط در آخرین لیست غیرقابل تغییر ارسالشده به downstream وجود دارد.
میتوانید اولین مثال ساده را برای دیدن این Gatherer در عمل بنویسید.
اجرای مثال قبلی خروجی زیر را چاپ میکند.
result = [[1, 2, 3], [2, 3, 4], [3, 4, 5], [4, 5, 6], [5, 6, 7], [6, 7, 8], [7, 8, 9]]
توجه کنید برای این Gatherer، تمام پنجرههای تولیدشده اندازهی یکسانی دارند، تا زمانی که upstream عناصر بیشتری از اندازهی پنجرهی انتخابی شما تولید کند.
چون عناصر استریم gather شده خود لیست هستند، میتوانید آنها را بیشتر با Stream API پردازش کنید. بیایید این کار را با ایجاد gathererای که میانگین هر پنجره را محاسبه میکند، با ترکیب دو gatherer انجام دهیم.
اجرای کد قبلی خروجی زیر را چاپ میکند.
result = [2.0, 3.0, 4.0, 5.0, 6.0, 7.0, 8.0]
این محتوا کاملا رایگان توسط تیم کدلپر ترجمه شده و در اختیار شما کاربران عزیز قرار گرفته است، هر گونه کپی برداری برای مقاصد غیر رایگان و بدون ذکر منبع، مورد پیگیری قانونی قرار میگیرد.
ترجمه شده از منبع: https://dev.java/learn/