Kotlin: شرح Flow في كوتلن
آخر تحديث: 2026-08-26
Flow هو امتداد التدفق للكوروتينات — التدفقات الباردة تُرسل عند الطلب، الضغط العكسي مدمج، والعوامل قابلة للتركيب. يستخدم تشارلي Flow لبناء خط معالجة طلبات في الوقت الفعلي.
1. ما ستتعلمه
- التدفق البارد مقابل الساخن:
Flowيُرسل عند الطلب - العمليات الوسيطة:
map/filter/transform/take/debounce - العمليات النهائية:
collect/toList/first/single - معالجة الضغط العكسي:
buffer/conflate/collectLatest - تشارلي في العمل: خط تدفق أحداث الطلبات
2. قصة مهندس حقيقية
(1) نقطة الألم: عنق زجاجة معالجة الطلبات في الوقت الفعلي
يحتاج OrderProcessor لتشارلي معالجة 5,000 حدث طلب في الثانية في الوقت الفعلي. المعالجة الدفعية بـ List ذات زمن انتقال عالٍ، وتداخل الاستدعاءات غير قابل للصيانة، و RxJava منحنى تعلم حاد.
(2) حل Flow
KOTLIN
// Flow: تدفق بارد، إرسال عند الطلب، ضغط عكسي مدمج
fun orderEventStream(): Flow<OrderEvent> = flow {
while (true) {
val event = pollOrderEvent() // استطلاع غير مانع
emit(event) // إرسال للأسفل
}
}
// تركيب خط أنابيب: تصفية -> تحويل -> جمع
orderEventStream()
.filter { it.total > 100 }
.map { enrich(it) }
.collect { dispatch(it) }
Flow هو تدفق أصلي للكوروتينات — إرسال بارد عند الطلب، ضغط عكسي مدمج، عوامل قابلة للتركيب، منحنى تعلم أقل بكثير من RxJava.
3. أساسيات Flow
(1) إنشاء التدفقات
KOTLIN
// 1. باني flow
fun numberFlow(): Flow<Int> = flow {
for (i in 1..5) {
emit(i) // إرسال قيمة للجامع
delay(100) // محاكاة عمل
}
}
// 2. flowOf: من قيم ثابتة
val orderFlow = flowOf(
Order("ORD-001", 299.99),
Order("ORD-002", 1_500.00)
)
// 3. asFlow: من Iterable/Sequence
val orderList = listOf(Order("ORD-001", 299.99))
val orderFlow2 = orderList.asFlow()
// 4. channelFlow: للإرسال المتزامن المعقد
fun concurrentFlow(): Flow<Int> = channelFlow {
launch { send(1) }
launch { send(2) }
}
(2) سلوك التدفق البارد
KOTLIN
// تدفق بارد: كل جامع يثير إرسالاً جديداً
val coldFlow = flow {
println("Emitting...") // هذا يعمل في كل مرة يُستدعى فيها collect
emit(1)
emit(2)
}
coldFlow.collect { println("Collector 1: $it") } // يرسل 1, 2
coldFlow.collect { println("Collector 2: $it") } // يرسل 1, 2 مرة أخرى
(3) التدفق البارد مقابل الساخن
| البُعد | التدفق البارد (Flow) | التدفق الساخن (SharedFlow/StateFlow) |
|---|---|---|
| المُثير | يرسل فقط عند الجمع | يعمل بشكل مستقل عن الجوامع |
| البيانات | كل جامع يحصل على البيانات الكاملة | الجوامع الجدد يرون البيانات اللاحقة فقط |
| التخزين المؤقت | لا شيء | قابل للتهيئة (replay) |
| دورة الحياة | تتبع الجامع | مستقلة |
| التشبيه | فيديو عند الطلب | بث مباشر |
4. العمليات الوسيطة
العمليات الوسيطة باردة — لا تثير الإرسال، بل تحدد فقط منطق التحويل.
(1) map / filter / transform
KOTLIN
// map: تحويل كل قيمة
val totals = orderFlow.map { it.total }
// filter: الاحتفاظ بالقيم المطابقة
val highValue = orderFlow.filter { it.total > 1_000 }
// transform: أكثر مرونة من map (يمكن إرسال قيم متعددة)
val expanded = orderFlow.transform { order ->
emit("Order: ${order.id}")
emit("Total: \$${order.total} USD")
}
(2) take / drop / debounce
KOTLIN
// take: أخذ أول N قيمة
val first3 = orderFlow.take(3)
// drop: تخطي أول N قيمة
val after5 = orderFlow.drop(5)
// debounce: إرسال فقط بعد فترة صمت
val debounced = searchFlow.debounce(300) // انتظر 300ms بدون قيم جديدة
(3) فئات العمليات الوسيطة
| الفئة | العوامل | الوظيفة |
|---|---|---|
| تحويل | map / transform |
تحويل لكل قيمة |
| تصفية | filter / filterNot / filterIsInstance |
تصفية شرطية |
| تسطيح | flatMapConcat / flatMapMerge / flatMapLatest |
تسطيح تدفقات متداخلة |
| اقتطاع | take / drop / takeWhile |
أخذ/تخطي أول N |
| توقيت | debounce / sample / throttleFirst |
تحكم زمني |
| تركيب | zip / combine |
دمج تدفقات متعددة |
5. العمليات النهائية
العمليات النهائية تثير الإرسال والجمع الفعلي.
(1) collect / toList / first
KOTLIN
// collect: استهلاك كل قيمة
orderFlow.collect { order -> println(order.id) }
// toList: جمع في قائمة
val allOrders = orderFlow.toList()
// first: الحصول على أول قيمة
val firstOrder = orderFlow.first()
// single: توقع قيمة واحدة بالضبط
val onlyOrder = orderFlow.single()
(2) فئات العمليات النهائية
| العملية | نوع الإرجاع | الوصف |
|---|---|---|
collect |
Unit | استهلاك كل قيمة |
toList / toSet |
List / Set |
جمع في مجموعة |
first / firstOrNull |
T / T? |
الحصول على أول قيمة |
single / singleOrNull |
T / T? |
الحصول على قيمة فريدة |
count |
Int | عد العناصر |
fold / reduce |
R / T |
تراكم |
6. معالجة الضغط العكسي
(1) مشكلة الضغط العكسي
KOTLIN
// المنتج أسرع من المستهلك
val fastFlow = flow {
repeat(10) {
emit(it) // سريع: 1ms لكل إرسال
delay(1)
}
}
// مستهلك بطيء يسبب ضغطاً عكسياً
fastFlow.collect {
delay(100) // بطيء: 100ms لكل معالجة
println("Processed: $it")
}
// افتراضياً: المنتج ينتظر المستهلك (تعليق)
(2) استراتيجيات الضغط العكسي
KOTLIN
// buffer: تشغيل المنتج والمستهلك بالتوازي
fastFlow.buffer(10) // تخزين مؤقت حتى 10 قيم
.collect { delay(100); println("Buffered: $it") }
// conflate: تخطي القيم الوسيطة، الاحتفاظ بالأحدث
fastFlow.conflate()
.collect { delay(100); println("Latest: $it") }
// collectLatest: إلغاء الجمع السابق عند وصول قيمة جديدة
fastFlow.collectLatest { value ->
delay(100)
println("Latest processed: $value")
}
(3) مقارنة استراتيجيات الضغط العكسي
| الاستراتيجية | السلوك | فقدان البيانات | حالة الاستخدام |
|---|---|---|---|
| افتراضي | المنتج ينتظر المستهلك | لا فقدان | ضمان الاكتمال |
buffer |
تخزين مؤقت لـ N قيمة | توقف عند امتلاء المخزن | ارتفاعات إنتاج قصيرة |
conflate |
تخطي الوسيطة، الاحتفاظ بالأحدث | القيم الوسيطة مفقودة | عرض في الوقت الفعلي |
collectLatest |
إلغاء المعالجة القديمة، معالجة الأحدث | المعالجة القديمة مُقاطعة | اقتراحات البحث |
7. خط أنابيب التدفق البارد
flowchart LR
A[المصدر<br/>flow emit] --> B[filter<br/>total > 100]
B --> C[map<br/>enrich]
C --> D[debounce<br/>300ms]
D --> E[collect<br/>dispatch]
8. مثال كامل: تدفق أحداث OrderProcessor
KOTLIN
// ============================================
// OrderProcessor - تدفق الأحداث مع Flow
// الميزة: معالجة أحداث الطلبات في الوقت الفعلي
// ============================================
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
data class Order(val id: String, val total: Double, val status: String, val customer: String)
sealed class OrderEvent {
data class Created(val order: Order) : OrderEvent()
data class Updated(val orderId: String, val newStatus: String) : OrderEvent()
data class Cancelled(val orderId: String, val reason: String) : OrderEvent()
}
// محاكاة مصدر الأحداث
fun orderEventStream(): Flow<OrderEvent> = flow {
val events = listOf(
OrderEvent.Created(Order("ORD-001", 299.99, "PENDING", "Alice")),
OrderEvent.Created(Order("ORD-002", 15_000.00, "PENDING", "Bob")),
OrderEvent.Updated("ORD-001", "CONFIRMED"),
OrderEvent.Created(Order("ORD-003", 2_500.00, "PENDING", "Charlie")),
OrderEvent.Cancelled("ORD-003", "Out of stock"),
OrderEvent.Updated("ORD-002", "CONFIRMED"),
OrderEvent.Created(Order("ORD-004", 8_900.00, "PENDING", "Bob")),
OrderEvent.Updated("ORD-002", "SHIPPED")
)
events.forEach { event ->
emit(event)
delay(100)
}
}
// إثراء الطلبات المنشأة
fun Order.enrich() = copy(status = "ENRICHED")
fun main() = runBlocking {
println("=== Order Event Stream ===")
// خط أنابيب: تصفية المنشأة -> إثراء -> جمع
orderEventStream()
.onStart { println("Stream started") }
.onEach { println(" Received: $it") }
.filter { it is OrderEvent.Created }
.map { (it as OrderEvent.Created).order }
.filter { it.total > 1_000 }
.map { it.enrich() }
.onEach { println(" >>> Enriched high-value: ${it.id} (\$${it.total} USD)") }
.onCompletion { println("Stream completed") }
.collect()
// خط أنابيب منفصل: تحديثات الحالة
println("\n=== Status Updates ===")
orderEventStream()
.filter { it is OrderEvent.Updated }
.map { it as OrderEvent.Updated }
.collect { event ->
println(" ${event.orderId} -> ${event.newStatus}")
}
// تجميع: عد الأحداث حسب النوع
println("\n=== Event Counts ===")
val counts = orderEventStream()
.groupBy { it::class.simpleName ?: "Unknown" }
.map { (type, flow) ->
type to flow.count()
}
.toList()
.associate { it }
counts.forEach { (type, count) ->
println(" $type: $count events")
}
}
المخرجات:
TEXT
📖 للعرض فقط
=== Order Event Stream ===
Stream started
Received: Created(order=Order(ORD-001, 299.99, PENDING, Alice))
Received: Created(order=Order(ORD-002, 15000.0, PENDING, Bob))
>>> Enriched high-value: ORD-002 ($15000.0 USD)
Received: Updated(orderId=ORD-001, newStatus=CONFIRMED)
Received: Created(order=Order(ORD-003, 2500.0, PENDING, Charlie))
>>> Enriched high-value: ORD-003 ($2500.0 USD)
Received: Cancelled(orderId=ORD-003, reason=Out of stock)
Received: Updated(orderId=ORD-002, newStatus=CONFIRMED)
Received: Created(order=Order(ORD-004, 8900.0, PENDING, Bob))
>>> Enriched high-value: ORD-004 ($8900.0 USD)
Received: Updated(orderId=ORD-002, newStatus=SHIPPED)
Stream completed
=== Status Updates ===
ORD-001 -> CONFIRMED
ORD-002 -> CONFIRMED
ORD-002 -> SHIPPED
=== Event Counts ===
Created: 4 events
Updated: 3 events
Cancelled: 1 events
9. أمثلة عملية سريعة
▶ مثال: Flow أساسي و collect
KOTLIN
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
fun main() = runBlocking {
// إنشاء Flow بسيط
val simpleFlow: Flow<Int> = flow {
for (i in 1..5) {
delay(100)
emit(i) // ينبعث قيمة لكل تكرار
}
}
// جمع القيم
println("--- Basic flow ---")
simpleFlow.collect { value ->
println("Got: $value")
}
// flowOf - Flow من قائمة ثابتة
val listFlow = flowOf(1, 2, 3, 4, 5)
val sum = listFlow.reduce { acc, value -> acc + value }
println("Sum: $sum")
}
**المخرجات:
TEXT
📖 للعرض فقط
--- Basic flow ---
Got: 1
Got: 2
Got: 3
Got: 4
Got: 5
Sum: 15
▶ مثال: عوامل Flow التحويلية
KOTLIN
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
fun main() = runBlocking {
val numbersFlow = (1..10).asFlow()
// map - تحويل كل قيمة
val doubledFlow = numbersFlow.map { it * 2 }
println("Doubled: ${doubledFlow.toList()}")
// filter - تصفية
val evensFlow = numbersFlow.filter { it % 2 == 0 }
println("Evens: ${evensFlow.toList()}")
// transform - تحويل عام
val transformedFlow = numbersFlow.transform { value ->
emit(value)
emit(value * value)
}
println("Transformed: ${transformedFlow.toList()}")
// chain من العوامل
val result = numbersFlow
.filter { it > 5 }
.map { it * 10 }
.take(3)
.toList()
println("Chain result: $result")
}
**المخرجات:
TEXT
📖 للعرض فقط
Doubled: [2, 4, 6, 8, 10, 12, 14, 16, 18, 20]
Evens: [2, 4, 6, 8, 10]
Transformed: [1, 1, 2, 4, 3, 9, 4, 16, 5, 25, 6, 36, 7, 49, 8, 64, 9, 81, 10, 100]
Chain result: [60, 70, 80]
▶ مثال: StateFlow و SharedFlow
KOTLIN
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
class TimerViewModel {
private val _seconds = MutableStateFlow(0)
val seconds: StateFlow<Int> = _seconds.asStateFlow()
fun start() {
// في الإنتاج: viewModelScope.launch { ... }
kotlinx.coroutines.GlobalScope.launch {
while (true) {
delay(1000)
_seconds.value++
}
}
}
}
fun main() = runBlocking {
val viewModel = TimerViewModel()
viewModel.start()
// جمع من StateFlow
val job = launch {
viewModel.seconds.collect { value ->
println("Seconds: $value")
}
}
delay(3500)
job.cancelAndJoin()
// SharedFlow - بث متعدد المشتركين
val sharedFlow = MutableSharedFlow<String>(replay = 1)
launch { sharedFlow.emit("Event A") }
launch { sharedFlow.emit("Event B") }
delay(100)
sharedFlow.collect { println("Got: $it") }
}
**المخرجات:
TEXT
📖 للعرض فقط
Seconds: 0
Seconds: 1
Seconds: 2
Seconds: 3
Got: Event B
▶ مثال: flatMap و merge
KOTLIN
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
fun main() = runBlocking {
// flatMapConcat - ينفذ بالتسلسل
val flatConcat = flowOf(1, 2, 3).flatMapConcat { value ->
flow {
delay(100)
emit("$value-a")
delay(100)
emit("$value-b")
}
}
println("flatMapConcat: ${flatConcat.toList()}")
// flatMapMerge - ينفذ بالتوازي
val flatMerge = flowOf(1, 2, 3).flatMapMerge { value ->
flow {
delay(100)
emit("$value-a")
delay(100)
emit("$value-b")
}
}
println("flatMapMerge: ${flatMerge.toList()}")
// flatMapLatest - يلغي التدفق السابق
val flatLatest = flowOf(1, 2, 3).flatMapLatest { value ->
flow {
delay(100)
emit("$value-a")
delay(100)
emit("$value-b")
}
}
println("flatMapLatest: ${flatLatest.toList()}")
}
**المخرجات:
TEXT
📖 للعرض فقط
flatMapConcat: [1-a, 1-b, 2-a, 2-b, 3-a, 3-b]
flatMapMerge: [1-a, 2-a, 3-a, 1-b, 2-b, 3-b]
flatMapLatest: [3-a, 3-b]
▶ مثال: عوامل طرفية (Terminal Operators)
KOTLIN
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
fun main() = runBlocking {
val numbers = (1..10).asFlow()
// reduce - تجميع
val sum = numbers.reduce { acc, value -> acc + value }
println("Sum: $sum")
// fold - مع قيمة ابتدائية
val product = numbers.fold(1) { acc, value -> acc * value }
println("Product: $product")
// first - أول قيمة
val firstEven = numbers.first { it % 2 == 0 }
println("First even: $firstEven")
// firstOrNull - آمن عند عدم العثور
val firstBig = numbers.firstOrNull { it > 100 }
println("First > 100: $firstBig")
// toList و toSet
val list = numbers.toList()
val set = numbers.toSet()
println("List: $list, Set: $set")
// count
val count = numbers.count { it > 5 }
println("Count > 5: $count")
// sumOf
val sumSquares = numbers.sumOf { it * it }
println("Sum of squares: $sumSquares")
}
**المخرجات:
TEXT
📖 للعرض فقط
Sum: 55
Product: 3628800
First even: 2
First > 100: null
List: [1, 2, 3, 4, 5, 6, 7, 8, 9, 10], Set: [1, 2, 3, 4, 5, 6, 7, 8, 9, 10]
Count > 5: 5
Sum of squares: 385
▶ مثال: buffer و conflate
KOTLIN
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
fun main() = runBlocking {
// بدون buffer - بطيء لأن كل emit يجب انتظار collect
val time1 = System.currentTimeMillis()
flow {
repeat(5) { i ->
delay(100) // محاكاة إنتاج بطيء
emit(i)
}
}.collect { value ->
delay(200) // محاكاة معالجة بطيئة
println("Slow: $value")
}
println("Slow time: ${System.currentTimeMillis() - time1}ms")
// مع buffer(2) - فصل الإنتاج عن الاستهلاك
val time2 = System.currentTimeMillis()
flow {
repeat(5) { i ->
delay(100)
emit(i)
}
}.buffer(2).collect { value ->
delay(200)
println("Buffered: $value")
}
println("Buffered time: ${System.currentTimeMillis() - time2}ms")
// conflate - يحتفظ فقط بآخر قيمة إذا كان المستهلك بطيئاً
val time3 = System.currentTimeMillis()
flow {
repeat(5) { i ->
delay(100)
emit(i)
}
}.conflate().collect { value ->
delay(250)
println("Conflated: $value")
}
println("Conflated time: ${System.currentTimeMillis() - time3}ms")
}
**المخرجات:
TEXT
📖 للعرض فقط
Slow: 0
Slow: 1
...
Slow time: ~1500ms
Buffered: 0
Buffered: 1
...
Buffered time: ~1100ms
Conflated: 0
Conflated: 4
Conflated time: ~650ms
▶ مثال: معالجة الأخطاء في Flow
KOTLIN
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
fun main() = runBlocking {
// Flow مع استثناء
val errorFlow = flow {
emit(1)
emit(2)
throw RuntimeException("Something failed!")
}
// catch - التقاط الاستثناءات
val recoveredFlow = errorFlow.catch { e ->
println("Caught: ${e.message}")
emit(-1) // قيمة استرداد
}
println("Recovered: ${recoveredFlow.toList()}")
// retry - إعادة المحاولة
var attempts = 0
val retriedFlow = flow {
attempts++
println("Attempt $attempts")
if (attempts < 3) {
throw RuntimeException("Failed attempt $attempts")
}
emit("Success after $attempts attempts")
}.retry(2) { e ->
println("Retrying due to: ${e.message}")
true // أعد المحاولة
}
println("Retry result: ${retriedFlow.toList()}")
// retryWhen - إعادة المحاولة الشرطية
val retriedWhenFlow = flow {
emit("OK")
}.retryWhen { _, attempt ->
attempt < 3 // حاول 3 مرات
}
}
**المخرجات:
TEXT
📖 للعرض فقط
Caught: Something failed!
Recovered: [1, 2, -1]
Attempt 1
Retrying due to: Failed attempt 1
Attempt 2
Retrying due to: Failed attempt 2
Attempt 3
Retry result: [Success after 3 attempts]
❓ أسئلة شائعة
س ما الفرق بين Flow و Sequence؟
ج Flow غير متزامن (يدعم suspend)، Sequence متزامن. Flow يناسب مصادر البيانات غير المتزامنة (شبكة/قاعدة بيانات)؛ Sequence يناسب معالجة البيانات الكبيرة المتزامنة.
س ما الفرق بين Flow و Observable في RxJava؟
ج Flow أصلي للكوروتينات، أولوية التدفق البارد، منحنى تعلم لطيف. RxJava أكثر نضجاً لكن أكثر تعقيداً. اختر Flow للمشاريع الجديدة؛ تعايش مع RxJava في المشاريع الموجودة.
س متى يصبح التدفق البارد تدفقاً ساخناً؟
ج استخدم
shareIn لتحويل Flow بارد إلى SharedFlow، أو stateIn لـ StateFlow. هذا يتيح لجوامع متعددي مشاركة نفس مصدر البيانات.س كيف أختار بين flatMapConcat و flatMapMerge و flatMapLatest؟
ج
flatMapConcat يعالج بالتسلسل (افتراضي)، flatMapMerge يعالج بالتوازي (إنتاجية أعلى)، flatMapLatest يلغي التدفقات القديمة (سيناريوهات الوقت الفعلي مثل البحث).س كيف يتعامل Flow مع الضغط العكسي افتراضياً؟
ج السلوك الافتراضي هو "المنتج ينتظر المستهلك" —
emit دالة suspend؛ عندما يكون المستهلك بطيئاً، يُعلّق المنتج تلقائياً. هذا هو الافتراضي الأكثر أماناً.س هل collect دالة suspend؟
ج نعم.
collect دالة suspend، قابلة للاستدعاء فقط من كوروتينات أو دوال suspend. إنها تُعلّق حتى يكتمل Flow.📖 ملخص
- Flow تدفق بارد: كل جامع يثير إرسالاً مستقلاً، لا مشاركة للبيانات
- العمليات الوسيطة (map/filter/transform) باردة — لا تثير الإرسال
- العمليات النهائية (collect/toList/first) تثير تدفق البيانات الفعلي
- الضغط العكسي الافتراضي: المنتج ينتظر المستهلك؛ عدّل بـ buffer/conflate/collectLatest
- صياغة عوامل Flow متسقة مع عوامل المجموعات — منحنى تعلم منخفض
onStart/onEach/onCompletionهي خطافات دورة حياة التدفق
📝 تمارين
- مبتدئ (⭐): أنشئ Flow يرسل 1-10 باستخدام
flow { }، ثم صفِّ واطبع الأرقام الزوجية. تلميح:flow { for (i in 1..10) emit(i) }.filter { } - متوسط (⭐⭐): حاكِ تدفق أحداث طلبات مع Flow، واستخدم
filter+map+collectلمعالجة الطلبات بمبالغ > 1000. تلميح: راجع القسم 8 - متقدم (⭐⭐⭐): نفّذ تدفق بحث مع
debounce: حاكِ إدخال المستخدم، أثر البحث فقط بعد 300ms بدون إدخال جديد. تلميح:debounce(300)