बैकएंड इंजीनियरिंग और कार्यात्मक प्रोग्रामिंग के लिए Scala · पाठ

Streams को बदलना

प्रवाहित डेटा को मैप और फ़िल्टर करें।

पाठ 2, कुल 4 में से13 चरण

Streams को बदलना, CoddyKit पर बैकएंड इंजीनियरिंग और कार्यात्मक प्रोग्रामिंग के लिए Scala का एक निःशुल्क पाठ है। यह 4 में से 2वाँ पाठ है। आप नीचे पूरा पाठ निःशुल्क पढ़ सकते हैं—फिर अंतर्निहित कोड संपादक और 24/7 एआई ट्यूटर के साथ ब्राउज़र में इसका व्यावहारिक अभ्यास कर सकते हैं। यह बैकएंड इंजीनियरिंग और कार्यात्मक प्रोग्रामिंग के लिए Scala सीखने के मार्ग का हिस्सा है और आपकी प्रगति वेब तथा CoddyKit ऐप पर सिंक होती रहती है। बैकएंड इंजीनियरिंग और कार्यात्मक प्रोग्रामिंग के लिए Scala पाठ्यक्रम में कुल 4 पाठ शामिल हैं।

रूपांतरण के रूप में ऑपरेटर

Akka Streams, Source और Flow पर ऑपरेटरों का एक समृद्ध सेट उपलब्ध कराता है। ये Scala के कलेक्शन API जैसे हैं, लेकिन असिंक्रोनस रूप से चलते हैं और बैकप्रेशर का सम्मान करते हैं।

हर ऑपरेटर एक नया ब्लूप्रिंट लौटाता है, इसलिए स्ट्रीम चलने से पहले ही रूपांतरणों को घोषणात्मक ढंग से संयोजित किया जाता है।

map और filter

map हर तत्व पर एक सिंक्रोनस फ़ंक्शन लागू करता है; filter उन तत्वों को हटा देता है जो किसी प्रेडिकेट को संतुष्ट नहीं करते। ये तत्व-दर-तत्व रूपांतरण के सबसे उपयोगी ऑपरेटर हैं।

दोनों क्रम को बनाए रखते हैं और पूर्णता तथा विफलता को डाउनस्ट्रीम तक पहुँचाते हैं।

val flow =
  Flow[Int]
    .filter(_ % 2 == 0)
    .map(n => n * n)

एक से अनेक के लिए mapConcat

जब एक इनपुट से कई आउटपुट बनाने हों, तब mapConcat का उपयोग करें। यह ऐसा फ़ंक्शन लेता है जो iterable लौटाता है और परिणामों को स्ट्रीम में समतल कर देता है।

खाली कलेक्शन लौटाने पर वह तत्व प्रभावी रूप से हटा दिया जाता है।

val explode: Flow[String, String, akka.NotUsed] =
  Flow[String].mapConcat(line => line.split(",").toList)

val words = Source(List("a,b", "c,d,e"))
  .via(explode)

grouped और sliding

grouped(n) लगातार आने वाले तत्वों को अधिकतम n आइटम वाले Seq में समूहित करता है। यह बड़े पैमाने पर डेटाबेस लेखन के लिए उपयोगी है। sliding(n) एक-दूसरे पर आंशिक रूप से चढ़ने वाली विंडो उत्सर्जित करता है।

बैच बनाने से I/O-प्रधान पाइपलाइनों में प्रति-तत्व अतिरिक्त लागत कम होती है।

val batches: Source[Seq[Int], akka.NotUsed] =
  Source(1 to 1000).grouped(100)

val windows =
  Source(1 to 10).sliding(3, step = 1)

scan और fold

scan हर तत्व के बाद चल रहे संचायक का मान उत्सर्जित करता है, जिससे बदलती हुई स्थिति की स्ट्रीम मिलती है। fold अपस्ट्रीम पूरा होने के बाद केवल अंतिम संचित मान एक बार उत्सर्जित करता है।

लाइव काउंटर के लिए scan और अंतिम समेकित परिणामों के लिए fold का उपयोग करें।

val running =
  Source(1 to 5).scan(0)(_ + _) // 0,1,3,6,10,15

val total =
  Source(1 to 5).fold(0)(_ + _) // 15

असिंक्रोनस कार्य के लिए mapAsync

mapAsync(parallelism) ऐसा फ़ंक्शन कॉल करता है जो Future लौटाता है और परिणामों को क्रम में उत्सर्जित करता है। यह अधिकतम parallelism Futures को एक साथ चलाता है।

जहाँ क्रम महत्वपूर्ण हो, वहाँ डेटाबेस लुकअप या HTTP अनुरोध जैसे असिंक्रोनस कॉल के लिए इसका उपयोग करें।

import scala.concurrent.Future

val enriched =
  Flow[UserId]
    .mapAsync(parallelism = 4)(id => lookup(id))

def lookup(id: UserId): Future[User] = ???

mapAsyncUnordered

mapAsyncUnordered, mapAsync की तरह काम करता है, लेकिन हर परिणाम पूरा होते ही उसे उत्सर्जित करता है और इनपुट क्रम की अनदेखी करता है।

जब डाउनस्ट्रीम को क्रम की परवाह न हो, तब यह थ्रूपुट बढ़ा सकता है, क्योंकि धीमा Future अब तेज़ Futures को रोकता नहीं है।

val fast =
  Flow[UserId]
    .mapAsyncUnordered(parallelism = 8)(id => lookup(id))

statefulMapConcat के साथ स्थिति-आधारित रूपांतरण

ऐसे प्रति-तत्व रूपांतरणों के लिए जिन्हें स्थानीय परिवर्तनीय स्थिति चाहिए, statefulMapConcat हर मटेरियलाइज़ेशन के लिए नई स्थिति बनाता है और आउटपुट का iterable लौटाता है।

स्ट्रीम के अलग-अलग रन के बीच स्थिति साझा किए बिना काउंटर या बफ़र बनाए रखने का यह सुरक्षित तरीका है।

val withIndex: Flow[String, (Int, String), akka.NotUsed] =
  Flow[String].statefulMapConcat { () =>
    var i = 0
    elem => { i += 1; List((i, elem)) }
  }

समय-आधारित ऑपरेटर

स्ट्रीम सामग्री के साथ-साथ समय के आधार पर भी रूपांतरित की जा सकती हैं। throttle उत्सर्जन की दर सीमित करता है, groupedWithin आकार या बीते हुए समय के आधार पर बैच बनाता है, और takeWithin अवधि सीमित करता है।

बाहरी API की दर सीमित करने के लिए ये अत्यंत आवश्यक हैं।

import scala.concurrent.duration._

val limited =
  Source(1 to 1000)
    .throttle(10, 1.second)
    .groupedWithin(100, 500.millis)

रूपांतरणों में त्रुटियों का प्रबंधन

किसी ऑपरेटर के भीतर उत्पन्न अपवाद डिफ़ॉल्ट रूप से पूरी स्ट्रीम को विफल कर देता है। इसके बजाय, supervision strategy resume के माध्यम से खराब तत्व को हटा सकती है या चरण को restart कर सकती है।

Flow पर withAttributes के साथ यह रणनीति जोड़ें।

import akka.stream.{ActorAttributes, Supervision}

val safe =
  Flow[String].map(_.toInt)
    .withAttributes(
      ActorAttributes.supervisionStrategy(_ => Supervision.Resume))

Flows को संयोजित करना

छोटे Flows को via की सहायता से बड़े Flows में संयोजित किया जा सकता है, जिससे एकल पुनः उपयोग योग्य Flow बनता है। इससे हर रूपांतरण अपने उद्देश्य पर केंद्रित और स्वतंत्र रूप से जाँचा जा सकने वाला रहता है।

संयोजित Flow का इनपुट प्रकार पहले Flow का और आउटपुट प्रकार अंतिम Flow का होता है।

val parse  = Flow[String].map(_.toInt)
val square = Flow[Int].map(n => n * n)

val parseAndSquare: Flow[String, Int, akka.NotUsed] =
  parse.via(square)

त्वरित जाँच

असिंक्रोनस रूपांतरणों और उनके क्रम संबंधी आश्वासनों पर विचार करें।

पुनरावलोकन

आपने रूपांतरण ऑपरेटरों का अध्ययन किया: तत्व-दर-तत्व map/filter, एक-से-अनेक mapConcat, बैच बनाने के लिए grouped, scan और fold से संचयन, तथा mapAsync के माध्यम से असिंक्रोनस कार्य।

आपने स्थिति-आधारित रूपांतरण, throttle जैसे समय-आधारित ऑपरेटर, त्रुटियों के लिए supervision strategy, और via से Flows को संयोजित करना भी देखा। अगला विषय है: बैकप्रेशर इन चरणों को सुरक्षित कैसे रखता है।

शुरुआत निःशुल्क

एआई शिक्षक के साथ Scala सीखें — निःशुल्क

अपने ब्राउज़र में वास्तविक कोड लिखें और चलाएँ, चौबीसों घंटे एआई शिक्षक से तुरंत सहायता पाएँ, और वेब या ऐप पर वहीं से शुरू करें जहाँ आपने छोड़ा था।

पाठ्यक्रम
39
पाठ
143

अक्सर पूछे जाने वाले प्रश्न

क्या “Streams को बदलना” पाठ निःशुल्क है?

हाँ—“Streams को बदलना” का पूरा पाठ यहाँ वेब पर निःशुल्क पढ़ा जा सकता है। इंटरैक्टिव अभ्यास (अंतर्निहित कोड संपादक और 24/7 एआई ट्यूटर) करने और बैकएंड इंजीनियरिंग और कार्यात्मक प्रोग्रामिंग के लिए Scala पाठ्यक्रम का बाकी हिस्सा अनलॉक करने के लिए CoddyKit PRO लें। बैकएंड इंजीनियरिंग और कार्यात्मक प्रोग्रामिंग के लिए Scala पाठ्यक्रम में कुल 4 पाठ शामिल हैं।

“Streams को बदलना” में मैं क्या सीखूँगा?

प्रवाहित डेटा को मैप और फ़िल्टर करें। आप ब्राउज़र में सीधे चलाए जाने वाले व्यावहारिक कोड के साथ बैकएंड इंजीनियरिंग और कार्यात्मक प्रोग्रामिंग के लिए Scala का अभ्यास करते हैं, और पाठ पूरा करते समय 24/7 एआई ट्यूटर आपके प्रश्नों के उत्तर देता है।

क्या बैकएंड इंजीनियरिंग और कार्यात्मक प्रोग्रामिंग के लिए Scala शुरू करने के लिए मुझे किसी अनुभव की आवश्यकता है?

पहले के अनुभव की आवश्यकता नहीं है। CoddyKit पर बैकएंड इंजीनियरिंग और कार्यात्मक प्रोग्रामिंग के लिए Scala शुरुआती से लेकर उन्नत शिक्षार्थियों तक सभी के लिए व्यवस्थित किया गया है, इसलिए आप यहीं से या शुरुआत से सीखना शुरू कर सकते हैं और अपनी गति से आगे बढ़ सकते हैं। यह 4 में से 2वाँ पाठ है।

“Streams को बदलना” पाठ पूरा करने में कितना समय लगता है?

CoddyKit का अधिकांश पाठ लगभग 5–10 मिनट में पूरा हो जाता है। हर पाठ छोटा और संवादात्मक है, इसलिए आप लगातार प्रगति करते हैं और वेब या ऐप पर वहीं से सीखना जारी रख सकते हैं जहाँ आपने छोड़ा था।

क्या मैं इस बैकएंड इंजीनियरिंग और कार्यात्मक प्रोग्रामिंग के लिए Scala पाठ में कोड लिख और चला सकता हूँ?

हाँ। हर बैकएंड इंजीनियरिंग और कार्यात्मक प्रोग्रामिंग के लिए Scala पाठ में एक अंतर्निर्मित कोड संपादक शामिल है, जिससे आप सीधे अपने ब्राउज़र में वास्तविक कोड लिख और चला सकते हैं और तुरंत एआई प्रतिक्रिया पा सकते हैं—स्थानीय सेटअप की आवश्यकता नहीं है।

इस पाठ्यक्रम के सभी पाठ

  1. Source, Flow और Sink
  2. Streams को बदलना
  3. Backpressure
  4. Pipeline चलाना
← बैकएंड इंजीनियरिंग और कार्यात्मक प्रोग्रामिंग के लिए Scala पर वापस जाएँ