Kotlin: شرح Flow في كوتلن

آخر تحديث: 2026-08-26

Flow هو امتداد التدفق للكوروتينات — التدفقات الباردة تُرسل عند الطلب، الضغط العكسي مدمج، والعوامل قابلة للتركيب. يستخدم تشارلي Flow لبناء خط معالجة طلبات في الوقت الفعلي.

1. ما ستتعلمه


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. خط أنابيب التدفق البارد

100%
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.

📖 ملخص


📝 تمارين

  1. مبتدئ (⭐): أنشئ Flow يرسل 1-10 باستخدام flow { }، ثم صفِّ واطبع الأرقام الزوجية. تلميح: flow { for (i in 1..10) emit(i) }.filter { }
  2. متوسط (⭐⭐): حاكِ تدفق أحداث طلبات مع Flow، واستخدم filter + map + collect لمعالجة الطلبات بمبالغ > 1000. تلميح: راجع القسم 8
  3. متقدم (⭐⭐⭐): نفّذ تدفق بحث مع debounce: حاكِ إدخال المستخدم، أثر البحث فقط بعد 300ms بدون إدخال جديد. تلميح: debounce(300)

← السابق | التالي →

Web-Tutorial.com

فريق Web-Tutorial التقني

منصة دروس برمجية يديرها عدة مطورين. كل درس يتم كتابته ومراجعته بواسطة مطورين متخصصين في المجال. نعمل على ضمان دقة وموثوقية المحتوى — إذا لاحظت أي مشكلة، فيرجى إخبارنا.

100%