1 00:00:05,000 --> 00:00:11,000 Hello dear students, in this lesson we are going to learn what reactive programming is and how we can 2 00:00:11,000 --> 00:00:16,000 apply reactive programming practices using Java via flow API. 3 00:00:16,000 --> 00:00:20,000 We are going to start the lesson from learning general concepts. 4 00:00:20,000 --> 00:00:23,000 We'll learn what reactive programming approach is. 5 00:00:23,000 --> 00:00:29,000 We'll discuss what the most popular libraries for reactive programming in Java are. 6 00:00:29,000 --> 00:00:36,000 But in scope of this lesson, we will focus our attention on the flow API that is available in JDK for 7 00:00:36,000 --> 00:00:38,000 reactive programming. 8 00:00:38,000 --> 00:00:39,000 As we will go. 9 00:00:39,000 --> 00:00:45,000 We'll learn different terms like for example Backpressure aware concept. 10 00:00:45,000 --> 00:00:49,000 We'll learn key interfaces from flow API. 11 00:00:49,000 --> 00:00:51,000 Then we'll talk about use cases. 12 00:00:51,000 --> 00:00:57,000 And I will try to answer the question when it is recommended to use reactive programming approach. 13 00:00:57,000 --> 00:01:04,000 And at the end of the lesson we are going to review code examples in order to understand how reactive 14 00:01:04,000 --> 00:01:06,000 streams work on practice. 15 00:01:06,000 --> 00:01:08,000 Let's start our lesson. 16 00:01:09,000 --> 00:01:12,000 Let's start from learning of what reactive programming is. 17 00:01:13,000 --> 00:01:20,000 Reactive programming is a programming paradigm focused on handling asynchronous data streams and the 18 00:01:20,000 --> 00:01:22,000 propagation of changes. 19 00:01:22,000 --> 00:01:29,000 It provides a set of principles, patterns, and libraries that enable developers to build responsive 20 00:01:29,000 --> 00:01:36,000 and scalable systems, especially in the context of handling events, input output operations, and 21 00:01:36,000 --> 00:01:37,000 real time data. 22 00:01:38,000 --> 00:01:41,000 Let's review key characteristics of reactive programming. 23 00:01:41,000 --> 00:01:45,000 They include asynchronous and non-blocking. 24 00:01:46,000 --> 00:01:51,000 Reactive programming emphasizes asynchronous and non-blocking operations. 25 00:01:51,000 --> 00:01:58,000 It allows components to continue processing other tasks while waiting for asynchronous operations to 26 00:01:58,000 --> 00:02:02,000 complete, leading to more efficient resource utilization. 27 00:02:03,000 --> 00:02:04,000 Data streams. 28 00:02:04,000 --> 00:02:09,000 At the core of reactive programming is the concept of data streams. 29 00:02:09,000 --> 00:02:13,000 Sequences of events or values that occur over time. 30 00:02:13,000 --> 00:02:21,000 These streams can represent various types of data such as user inputs, sensor readings, or messages 31 00:02:21,000 --> 00:02:22,000 from a server. 32 00:02:23,000 --> 00:02:30,000 Event-driven reactive systems are often event driven, meaning they respond to events or changes in 33 00:02:30,000 --> 00:02:31,000 the system. 34 00:02:31,000 --> 00:02:37,000 Events can trigger reactions and the system adapts dynamically to these changes. 35 00:02:38,000 --> 00:02:39,000 Back pressure. 36 00:02:40,000 --> 00:02:43,000 Back pressure is a key concept in reactive programming. 37 00:02:43,000 --> 00:02:51,000 It refers to the ability of downstream component to signal its capacity to handle data to an upstream 38 00:02:51,000 --> 00:02:52,000 component. 39 00:02:52,000 --> 00:02:59,000 This helps prevent overwhelming the system with more data than it can process efficiently. 40 00:03:00,000 --> 00:03:02,000 Responsive and scalable. 41 00:03:02,000 --> 00:03:09,000 Reactive programming aims to create responsive applications that can handle varied workloads. 42 00:03:09,000 --> 00:03:17,000 By embracing asynchrony and scalability, reactive systems can efficiently adapt to changes in demand. 43 00:03:18,000 --> 00:03:20,000 Functional programming. 44 00:03:20,000 --> 00:03:25,000 Reactive programming often aligns with functional programming principles. 45 00:03:25,000 --> 00:03:33,000 It emphasizes immutability, pure functions, and the use of higher order functions, enabling developers 46 00:03:33,000 --> 00:03:36,000 to write more concise and maintainable code. 47 00:03:37,000 --> 00:03:39,000 Reactive extensions. 48 00:03:39,000 --> 00:03:47,000 Reactive extensions, often referred to as RCS, is a set of libraries that provide a uniform way to 49 00:03:47,000 --> 00:03:52,000 work with reactive programming concepts across multiple programming languages. 50 00:03:53,000 --> 00:03:58,000 RCS libraries exist for languages like Java, JavaScript, C, sharp, and more. 51 00:03:59,000 --> 00:04:05,000 Reactive programming has become popular in building modern applications, especially those that involve 52 00:04:05,000 --> 00:04:10,000 real time data, user interfaces, and distributed systems. 53 00:04:10,000 --> 00:04:18,000 Frameworks and libraries like Project Reactor, Erik's Java provide tools and abstractions to implement 54 00:04:18,000 --> 00:04:21,000 reactive programming in various programming languages. 55 00:04:22,000 --> 00:04:29,000 These tools simplify the handling of asynchronous and streaming data, making it easier for developers 56 00:04:29,000 --> 00:04:32,000 to build responsive and scalable applications. 57 00:04:33,000 --> 00:04:40,000 In Java, reactive programming approach may be implemented with different libraries, for example a 58 00:04:40,000 --> 00:04:42,000 Java reactor core. 59 00:04:42,000 --> 00:04:52,000 Reactive streams are socket martini reactor, but starting from Java version nine, you can use JDK 60 00:04:52,000 --> 00:04:55,000 to implement reactive programming approach. 61 00:04:55,000 --> 00:04:59,000 Flow API is available starting from Java nine. 62 00:04:59,000 --> 00:05:02,000 In this lesson we will focus on flow API. 63 00:05:03,000 --> 00:05:05,000 So let's learn what a flow API is. 64 00:05:06,000 --> 00:05:14,000 The flow API, also known as Reactive Streams, was introduced in Java nine to address challenges associated 65 00:05:14,000 --> 00:05:21,000 with asynchronous programming and handling streams of data in a more efficient and scalable way. 66 00:05:21,000 --> 00:05:28,000 Reactive streams provide a standard for asynchronous stream processing with non-blocking back pressure, 67 00:05:28,000 --> 00:05:33,000 enabling better control over the flow of data between components. 68 00:05:33,000 --> 00:05:40,000 Reactive streams define a set of interfaces, classes, and masses that standardize the exchange of 69 00:05:40,000 --> 00:05:46,000 asynchronous data streams between components in a non-blocking, back pressure aware manner. 70 00:05:47,000 --> 00:05:51,000 Now let me share with you key interfaces in the flow API. 71 00:05:51,000 --> 00:05:58,000 So the main interfaces include publisher, subscriber, subscription and processor. 72 00:05:58,000 --> 00:06:05,000 I will make just a brief overview of them and we will learn more during the review of real life example. 73 00:06:06,000 --> 00:06:13,000 Publisher amidst a potentially unbounded sequence of elements to one or more subscribers. 74 00:06:14,000 --> 00:06:22,000 Subscriber consumes elements from a publisher and signals when it can handle more elements. 75 00:06:22,000 --> 00:06:30,000 Subscription represents the link between a subscriber and a publisher, allowing the subscriber to request 76 00:06:30,000 --> 00:06:31,000 and cancel elements. 77 00:06:32,000 --> 00:06:33,000 Processor. 78 00:06:34,000 --> 00:06:43,000 The processor interface in the Reactive Streams API serves as a dual role component, acting both as 79 00:06:43,000 --> 00:06:49,000 a subscriber and a publisher, enabling custom data processing within a reactive stream. 80 00:06:50,000 --> 00:06:55,000 So let's try to answer the question why flow API is needed? 81 00:06:55,000 --> 00:07:02,000 Traditional stream processing in Java lacked a standardized way to handle asynchronous data streams, 82 00:07:02,000 --> 00:07:09,000 which could lead to issues like buffer overflow or excessive resource consumption. 83 00:07:09,000 --> 00:07:16,000 Reactive streams address these problems by introducing Backpressure, allowing subscribers to signal 84 00:07:16,000 --> 00:07:22,000 when they are ready to receive more data, preventing overwhelming downstream components. 85 00:07:22,000 --> 00:07:26,000 And in which cases it is recommended to use flow API. 86 00:07:27,000 --> 00:07:29,000 Asynchronous processing. 87 00:07:29,000 --> 00:07:36,000 Reactive streams are beneficial in scenarios where asynchronous processing is required, such as handling 88 00:07:36,000 --> 00:07:40,000 events, input output operations, or real time data. 89 00:07:41,000 --> 00:07:48,000 Backpressure handling when dealing with data streams, where the rate of data production is different 90 00:07:48,000 --> 00:07:50,000 from the rate of data consumption. 91 00:07:50,000 --> 00:07:56,000 Reactive streams help manage backpressure, preventing overwhelming downstream components. 92 00:07:57,000 --> 00:08:00,000 Scalable and responsive applications. 93 00:08:01,000 --> 00:08:02,000 Reactive programming was. 94 00:08:02,000 --> 00:08:09,000 Reactive Streams is well suited for building responsive and scalable applications, especially in scenarios 95 00:08:09,000 --> 00:08:13,000 like web services, where responsiveness is crucial. 96 00:08:14,000 --> 00:08:16,000 Use reactive streams when you need to. 97 00:08:17,000 --> 00:08:20,000 Handle asynchronous data streams. 98 00:08:20,000 --> 00:08:26,000 Implement Backpressure to control the flow of data and prevent resource exhaustion. 99 00:08:26,000 --> 00:08:34,000 Build scalable and responsive applications with non-blocking input output operations. 100 00:08:35,000 --> 00:08:41,000 Reactive streams have gained popularity in modern Java development and are widely used in frameworks 101 00:08:41,000 --> 00:08:46,000 like Project Reactor to build reactive and resilient systems. 102 00:08:46,000 --> 00:08:52,000 I believe we learned enough theory, and now we are ready to review examples that I prepared especially 103 00:08:52,000 --> 00:08:53,000 for this lesson. 104 00:08:54,000 --> 00:09:00,000 You can find the source code of examples that I am going to share with you in attachments to the lesson. 105 00:09:00,000 --> 00:09:06,000 Remember, Tim, that you are always welcome to ask questions in case something is not clear. 106 00:09:06,000 --> 00:09:11,000 Just post your question below the video and I will be happy to answer. 107 00:09:11,000 --> 00:09:17,000 As you can see, examples related to reactive streams are in flow API package. 108 00:09:18,000 --> 00:09:22,000 We're going to start review from the flow API example class. 109 00:09:22,000 --> 00:09:26,000 This is a class where I put client code and main method. 110 00:09:27,000 --> 00:09:30,000 So we will go together with the flow of execution. 111 00:09:30,000 --> 00:09:34,000 And I will explain you what is happening here line by line. 112 00:09:34,000 --> 00:09:39,000 Let me also run this file in order to have all console output printed. 113 00:09:39,000 --> 00:09:43,000 And you will be able to see the sequence of events that is happening here. 114 00:09:44,000 --> 00:09:51,000 Pay attention that today we will review types that are stored on the Java.util.concurrent package inside 115 00:09:51,000 --> 00:09:53,000 final class flow. 116 00:09:53,000 --> 00:10:00,000 There are such important interfaces as publisher, subscriber, subscription and processor. 117 00:10:00,000 --> 00:10:03,000 These are types that we are going to use today. 118 00:10:04,000 --> 00:10:07,000 I start example from creation of the submission publisher. 119 00:10:08,000 --> 00:10:15,000 The submission publisher is a built in implementation of the publisher interface in Java Flow API. 120 00:10:15,000 --> 00:10:23,000 It emits a potentially unbounded sequence of elements, in this case integers to one or more subscribers. 121 00:10:24,000 --> 00:10:29,000 The next thing I create instance of my subscriber. 122 00:10:29,000 --> 00:10:36,000 It is parametrized by integer type because in this example I plan to receive integer elements. 123 00:10:36,000 --> 00:10:42,000 To implement subscriber interface, you need to give implementation to all following methods. 124 00:10:42,000 --> 00:10:45,000 Let's review each of them one by one. 125 00:10:46,000 --> 00:10:48,000 Unsubscribe method. 126 00:10:48,000 --> 00:10:56,000 This method signals about the subscription sets the subscription object, and requests the first item 127 00:10:56,000 --> 00:10:59,000 to initialize the data flow with the provided subscription. 128 00:11:00,000 --> 00:11:07,000 Basically, with the request method, you create demand for elements that subscriber can process depending 129 00:11:07,000 --> 00:11:13,000 on the amount of requested elements, the same amount of times on next method will be invoked. 130 00:11:14,000 --> 00:11:16,000 On next method. 131 00:11:16,000 --> 00:11:21,000 This method processes in common integer items. 132 00:11:21,000 --> 00:11:28,000 In our example, I just double all in common integers and transform them to string just for the sake 133 00:11:28,000 --> 00:11:34,000 of example, because theoretically we can implement processor interface that I am going to explain a 134 00:11:34,000 --> 00:11:42,000 little bit later, and you can supply process data downstream in such way you can build processing chain 135 00:11:42,000 --> 00:11:43,000 and data flow. 136 00:11:44,000 --> 00:11:48,000 That's why I created submit method just for the sake of example. 137 00:11:48,000 --> 00:11:55,000 And just to show you that in submit method, you can put the logic of processed item submission downstream 138 00:11:55,000 --> 00:11:56,000 in simple words. 139 00:11:56,000 --> 00:12:03,000 Send processed item further down the processing chain depending on your business logic. 140 00:12:04,000 --> 00:12:12,000 Onerror method handles errors during processing by printing an error message, allowing the processor 141 00:12:12,000 --> 00:12:15,000 to respond appropriately to error conditions. 142 00:12:16,000 --> 00:12:20,000 On complete signals, the completion of processing. 143 00:12:20,000 --> 00:12:24,000 Printing a message to indicate the end of the data stream. 144 00:12:24,000 --> 00:12:31,000 If you build a longer processing chain, this is the method where you can notify downstream components 145 00:12:31,000 --> 00:12:36,000 that processing is completed, and in this method I cancel subscription. 146 00:12:37,000 --> 00:12:40,000 Let's get back to flow API example class. 147 00:12:41,000 --> 00:12:49,000 The next thing that I do here I subscribe my subscriber to publisher in the context of the Reactive 148 00:12:49,000 --> 00:12:50,000 Streams API. 149 00:12:50,000 --> 00:12:57,000 The relationship between a publisher and a subscriber is established through the subscribe method. 150 00:12:57,000 --> 00:13:04,000 However, it is important to note that it is the publisher who initiates the subscription process by 151 00:13:04,000 --> 00:13:08,000 calling the subscribe method on the publisher object. 152 00:13:08,000 --> 00:13:11,000 Now there will be a connection between them. 153 00:13:12,000 --> 00:13:19,000 Publisher will be aware about subscriber and when this will happen, then unsubscribe method will be 154 00:13:19,000 --> 00:13:24,000 triggered and I will initialize my subscription property inside. 155 00:13:25,000 --> 00:13:28,000 But wait, something should be wrong here. 156 00:13:28,000 --> 00:13:34,000 Where missing subscription object because we have publisher, we have subscriber. 157 00:13:34,000 --> 00:13:36,000 But where is subscription object? 158 00:13:36,000 --> 00:13:41,000 You can implement it by yourself by implementing subscription interface. 159 00:13:41,000 --> 00:13:43,000 But what is happening in this case? 160 00:13:44,000 --> 00:13:50,000 You can see that I decided to print subscription object to console in order to understand the type of 161 00:13:50,000 --> 00:13:51,000 this object. 162 00:13:52,000 --> 00:13:54,000 And we can see that we received here. 163 00:13:54,000 --> 00:13:56,000 Object of type buffered subscription. 164 00:13:57,000 --> 00:14:01,000 This is nested class inside submission publisher class. 165 00:14:01,000 --> 00:14:03,000 But who created this object? 166 00:14:03,000 --> 00:14:10,000 And when a buffered subscription is created when a subscriber subscribes to a submission publisher, 167 00:14:11,000 --> 00:14:17,000 this subscription object manages the communication between the publisher and the subscriber, facilitating 168 00:14:17,000 --> 00:14:19,000 the flow of data. 169 00:14:19,000 --> 00:14:21,000 The charm buffer. 170 00:14:21,000 --> 00:14:28,000 It suggests that this subscription might have some internal buffer to store items when the subscriber 171 00:14:28,000 --> 00:14:31,000 can't keep up with the rate of data production. 172 00:14:32,000 --> 00:14:36,000 Here is a breakdown of what happens during the subscription process. 173 00:14:37,000 --> 00:14:42,000 Publisher calls subscribe method and subscribe subscriber object. 174 00:14:42,000 --> 00:14:47,000 Submission publisher creates a buffered subscription object internally. 175 00:14:48,000 --> 00:14:51,000 This is happening inside subscribe method. 176 00:14:51,000 --> 00:14:57,000 If we would open the source code of subscribe method, we will see how instance of buffered subscription 177 00:14:57,000 --> 00:14:58,000 is created. 178 00:14:59,000 --> 00:15:05,000 Unsubscribe method is called on the subscriber with created buffered subscription. 179 00:15:05,000 --> 00:15:12,000 The unsubscribe method is then invoked on the subscriber, passing the buffered subscription as an argument. 180 00:15:12,000 --> 00:15:17,000 This allows the subscriber to request items and manage backpressure. 181 00:15:17,000 --> 00:15:24,000 The buffered aspect of the subscription typically comes into play when the subscriber is unable to keep 182 00:15:24,000 --> 00:15:27,000 up with the rate of data production. 183 00:15:27,000 --> 00:15:35,000 In such cases, the buffered subscription may buffer items until the subscriber signals that it can 184 00:15:35,000 --> 00:15:38,000 process more data by calling the request method. 185 00:15:39,000 --> 00:15:44,000 So after subscription is happened, then I set the timeout. 186 00:15:45,000 --> 00:15:49,000 We need this time out to make sure that subscription is created. 187 00:15:49,000 --> 00:15:56,000 But using timeout is very primitive way that I decided to use for the sake of example. 188 00:15:56,000 --> 00:16:03,000 In reality, ensuring that everyone is successfully subscribed in a reactive system can be a complicated 189 00:16:03,000 --> 00:16:09,000 task, and the approach might depend on the specific requirements of your application. 190 00:16:10,000 --> 00:16:14,000 Some points to consider using of countdown Lodge. 191 00:16:15,000 --> 00:16:17,000 Using of completablefuture. 192 00:16:17,000 --> 00:16:22,000 Using of timeout mechanism, including using it with completablefuture. 193 00:16:23,000 --> 00:16:26,000 In my course Java from zero to first job. 194 00:16:26,000 --> 00:16:33,000 I have multi-threading section where I explained in details API of countdown, Lodge, Completablefuture 195 00:16:33,000 --> 00:16:35,000 and other things. 196 00:16:35,000 --> 00:16:38,000 That's why I wouldn't stop too much on this right now. 197 00:16:39,000 --> 00:16:43,000 In console you can see subscriber subscribed. 198 00:16:44,000 --> 00:16:52,000 After the timeout, I declare a for loop with counter and I submit five integers using publisher object. 199 00:16:52,000 --> 00:17:00,000 I submit them one by one, and in console you can see that publisher submits one element and subscriber 200 00:17:00,000 --> 00:17:01,000 receives it. 201 00:17:01,000 --> 00:17:08,000 Taking into account that println method is synchronized, text and console not necessarily represents 202 00:17:08,000 --> 00:17:15,000 the exact sequence of events, but at least it can give us an understanding of what is going on. 203 00:17:15,000 --> 00:17:23,000 Subscriber after receiving integers, doubles, received elements and prints them to console two and 204 00:17:23,000 --> 00:17:30,000 when publisher is closed, subscriber got notified two and Oncomplete method is triggered. 205 00:17:31,000 --> 00:17:32,000 And that's it. 206 00:17:32,000 --> 00:17:39,000 Knowing this, you can build more complex data flows with multiple subscribers, different publishers. 207 00:17:39,000 --> 00:17:44,000 You can connect them between each other as your business logic would require. 208 00:17:45,000 --> 00:17:51,000 And the last but not the least thing that I would like to show you for today is processor interface. 209 00:17:52,000 --> 00:17:58,000 You can create processor class that would implement processor interface if you want to make it acting 210 00:17:58,000 --> 00:18:02,000 as both a publisher and a subscriber. 211 00:18:02,000 --> 00:18:06,000 This can be helpful when you want to build a processing chain. 212 00:18:06,000 --> 00:18:13,000 For example, you can receive some data from publisher, process it, and pass it further to the next 213 00:18:13,000 --> 00:18:14,000 subscriber. 214 00:18:15,000 --> 00:18:15,000 And like this. 215 00:18:15,000 --> 00:18:18,000 You can build complex data flows. 216 00:18:18,000 --> 00:18:26,000 So the processor interface in the Reactive Streams API is used when you need to create a component that 217 00:18:26,000 --> 00:18:29,000 acts as a both a publisher and a subscriber. 218 00:18:29,000 --> 00:18:36,000 It is particularly useful in scenarios where you want to transform or process the data flowing through 219 00:18:36,000 --> 00:18:40,000 a reactive stream before it reaches the downstream subscribers. 220 00:18:41,000 --> 00:18:46,000 Here are some situations where you might consider using the processor interface. 221 00:18:47,000 --> 00:18:50,000 Transforming or processing data. 222 00:18:50,000 --> 00:18:57,000 If you need to apply transformations, filtering or any other processing logic to the data before it 223 00:18:57,000 --> 00:19:03,000 is delivered to subscribers, the processor interface allows you to implement this logic. 224 00:19:04,000 --> 00:19:07,000 Connecting different parts of a reactive pipeline. 225 00:19:07,000 --> 00:19:13,000 When you have a reactive pipeline with multiple stages, and you want to connect these stages with a 226 00:19:13,000 --> 00:19:20,000 component that can both consume data from the upstream and emit processed data to the downstream. 227 00:19:21,000 --> 00:19:24,000 Implementing a custom flow control mechanism. 228 00:19:24,000 --> 00:19:31,000 If you need to implement a custom flow control mechanism between the publisher and subscribers, the 229 00:19:31,000 --> 00:19:38,000 processor allows you to manage backpressure and control the rate at which data is consumed. 230 00:19:38,000 --> 00:19:41,000 Aggregating or combining data streams. 231 00:19:41,000 --> 00:19:47,000 When you want to combine or aggregate data from multiple publishers before delivering it to downstream 232 00:19:47,000 --> 00:19:49,000 subscribers. 233 00:19:49,000 --> 00:19:52,000 A processor can be used to implement the necessary logic. 234 00:19:53,000 --> 00:20:00,000 If you check the source code of processor interface, then you will find that processor is a static 235 00:20:00,000 --> 00:20:05,000 interface that extends subscriber and publisher at the same time. 236 00:20:05,000 --> 00:20:10,000 In Java, an interface can extend multiple other interfaces. 237 00:20:10,000 --> 00:20:17,000 This feature allows an interface to inherit abstract masses from multiple parent interfaces, providing 238 00:20:17,000 --> 00:20:22,000 a way to express multiple types of behavior in a single interface. 239 00:20:22,000 --> 00:20:30,000 So feel free to use processor interface implementations when it will be needed in your business cases. 240 00:20:30,000 --> 00:20:34,000 That's all what I wanted to share with you today in this lesson. 241 00:20:34,000 --> 00:20:37,000 Let's recap what we have learned in the video. 242 00:20:38,000 --> 00:20:41,000 We learned reactive programming approach. 243 00:20:41,000 --> 00:20:45,000 Now you know the most popular libraries for reactive programming in Java. 244 00:20:46,000 --> 00:20:54,000 I explained you what flow API in JDK is, and we learned key interfaces and types in flow API. 245 00:20:54,000 --> 00:21:01,000 I explained when and in which cases it is recommended to use Reactive Programming, Reactive Streams 246 00:21:01,000 --> 00:21:08,000 flow API, and at the end of the lesson we reviewed code example in order to understand how it works 247 00:21:08,000 --> 00:21:09,000 in practice. 248 00:21:09,000 --> 00:21:12,000 That's all what I wanted to share with you today. 249 00:21:12,000 --> 00:21:15,000 Have a great day and see you in the next lesson.