Log.e(TAG, "onCompletion: ")

dismissLoading()

}

.catch {

Log.e(TAG, "catch: ")

mBinding.textView.text = “数据出错”

}

.collect {

Log.e(TAG, “collect: $it”)

mBinding.textView.text = “接收数据为:$it”

}

}

flowOn:线程切换(flow为IO,其他为Main)

lifecycleScope.launch {

Log.e(TAG, “flow:${Thread.currentThread()}”)

flow {

Log.e(TAG, “emit start:${Thread.currentThread()}”)

emit(“1”)

Log.e(TAG, “emit 1:${Thread.currentThread()}”)

emit(“2”)

Log.e(TAG, “emit 2:${Thread.currentThread()}”)

emit(“3”)

Log.e(TAG, “emit 3:${Thread.currentThread()}”)

emit(“4”)

Log.e(TAG, “emit 4:${Thread.currentThread()}”)

emit(“5”)

Log.e(TAG, “emit 5:${Thread.currentThread()}”)

emit(“6”)

Log.e(TAG, “emit 6:${Thread.currentThread()}”)

}

.flowOn(Dispatchers.IO)

.onStart {

Log.e(TAG, “onStart:${Thread.currentThread()}”)

showToast(“开始”)

}

.filter {

Log.e(TAG, “filter:${Thread.currentThread()}”)

it != “2”

}

.map {

Log.e(TAG, “map:${Thread.currentThread()}”)

“转换$it”

}

.transform<String,Int>{

Log.e(TAG, “transform1:${Thread.currentThread()}”)

emit( it.length)

Log.e(TAG, “transform2:${Thread.currentThread()}”)

}

// .zip(f1) { a, b ->

// Log.e(TAG, “zip:${Thread.currentThread()}”)

// “本流a:其他流a:其他流a:其他流b”

// }

.onCompletion {

Log.e(TAG, “onCompletion:${Thread.currentThread()}”)

showToast(“结束”)

}

.catch {

Log.e(TAG, “catch:${Thread.currentThread()}”)

showToast(“异常”)

}

.collect {

Log.e(TAG, “collect:${Thread.currentThread()}”)

Log.e(TAG, “collect:${it}”)

mBinding.textView.text = it.toString()

}

}

cancel:取消流

val job = lifecycleScope.launch {

Log.e(TAG, “flow:${Thread.currentThread()}”)

flow {

Log.e(TAG, “emit start:${Thread.currentThread()}”)

emit(“1”)

Log.e(TAG, “emit 1:${Thread.currentThread()}”)

emit(“2”)

Log.e(TAG, “emit 2:${Thread.currentThread()}”)

emit(“3”)

Log.e(TAG, “emit 3:${Thread.currentThread()}”)

emit(“4”)

Log.e(TAG, “emit 4:${Thread.currentThread()}”)

emit(“5”)

Log.e(TAG, “emit 5:${Thread.currentThread()}”)

emit(“6”)

Log.e(TAG, “emit 6:${Thread.currentThread()}”)

}

.flowOn(Dispatchers.IO)

.onStart {

Log.e(TAG, “onStart:${Thread.currentThread()}”)

showToast(“开始”)

}

.onCompletion {

Log.e(TAG, “onCompletion:${Thread.currentThread()}”)

showToast(“结束”)

}

.catch {

Log.e(TAG, “catch:${Thread.currentThread()}”)

showToast(“异常”)

}

.collect {

Log.e(TAG, “collect:${Thread.currentThread()}”)

Log.e(TAG, “collect:${it}”)

mBinding.textView.text = it.toString()

}

}

job.cancel()

filter :过滤操作符

lifecycleScope.launch {

//TODO List 转成 Flow

mList.asFlow()

.onEach {

delay(2000)

Log.e(TAG, “onEach: $it”)

}

.filter {

//TODO 数据过滤操作符

//只发送能被2整除的数据

Log.e(TAG, “filter: $it”)

it.toInt() % 2 == 0

}

.collect {

Log.e(TAG, “collect: $it”)

mBinding.textView.text = it

}

}

filterNot :过滤操作符

lifecycleScope.launch {

//TODO List 转成 Flow

mList.asFlow()

.onEach {

delay(2000)

Log.e(TAG, “onEach: $it”)

}

.filterNot {

//TODO 数据过滤操作符

//只发送不能被2整除的数据

Log.e(TAG, “filterNot: $it”)

it.toInt() % 2 == 0

}

.collect {

Log.e(TAG, “collect: $it”)

mBinding.textView.text = it

}

}

transform:转换操作符(需主动发送数据)

lifecycleScope.launch {

//TODO List 转成 Flow

// mList.asFlow()

flow {

//TODO 上游发射数据

Log.e(TAG, “emit1: start”)

emit(“1”)

Log.e(TAG, “emit2: start”)

emit(“2”)

Log.e(TAG, “emit3: start”)

emit(“3”)

Log.e(TAG, “emit: all end”)

}

.onEach {

delay(2000)

Log.e(TAG, “onEach: $it”)

}

.transform<String, Int> {

//TODO 转化操作符 转化完成后需主动发送数据

Log.e(TAG, “transform: $it”)

val value = it.toInt() + 100

emit(value)

}

.collect {

Log.e(TAG, “collect: $it”)

mBinding.textView.text = it.toString()

}

}

map:转换操作符

lifecycleScope.launch {

//TODO List 转成 Flow

mList.asFlow()

.onEach {

delay(2000)

Log.e(TAG, “onEach: $it”)

}

.map {

//TODO 数据类型转换操作 内部实现transform

Log.e(TAG, “map: $it”)

it.toInt() + 100

}

.collect {

Log.e(TAG, “collect: $it”)

mBinding.textView.text = it.toString()

}

}

take: 截取操作符

lifecycleScope.launch {

//TODO List 转成 Flow

mList.asFlow()

.onEach {

delay(2000)

Log.e(TAG, “onEach: $it”)

}

//TODO 截取操作符,截取N位发射数据

.take(2)

.collect {

Log.e(TAG, “collect: $it”)

mBinding.textView.text = it.toString()

}

}

buffer:背压操作符

lifecycleScope.launch {

//TODO List 转成 Flow

flow {

//TODO 上游发射数据

Log.e(TAG, “emit1: start”)

emit(1)

Log.e(TAG, “emit2: start”)

emit(2)

Log.e(TAG, “emit3: start”)

emit(3)

Log.e(TAG, “emit: all end”)

}

.onEach {

// delay(2000)

Log.e(TAG, “onEach: $it”)

}

//TODO 背压:.buffer() 先emit all 再collect

// .buffer()

//TODO 背压:.buffer(0) 先emit 1和2 -> collect 1和2 再emit 3 -> collect 3

.buffer(0)

.collect {

Log.e(TAG, “collect: $it”)

mBinding.textView.text = it.toString()

}

}

conflate

lifecycleScope.launch {

//TODO List 转成 Flow

flow {

//TODO 上游发射数据

for(i in 1…30) {

delay(100)

emit(i)

}

}

//TODO 仅保留最新的值,无论上游发射多少数据,下游只会接收最新值

.conflate()

.onEach {

delay(2000)

Log.e(TAG, “onEach: $it”)

}

.collect {

Log.e(TAG, “collect: $it”)

mBinding.textView.text = it.toString()

}

}

zip:流合并操作符

val flowOther = (101…110).asFlow()

lifecycleScope.launch {

//TODO List 转成 Flow

mList.asFlow()

.onEach {

delay(2000)

Log.e(TAG, “onEach: $it”)

}

//TODO 合并其他流,本流发射数 < 其他流的发射数时,合并完的次数为本流的次数

.zip(flowOther){a,b->

“本流a:其他流a:其他流a:其他流b”

}

.collect {

Log.e(TAG, “collect: $it”)

mBinding.textView.text = it.toString()

}

}

combine:流组合操作符

val flowOther = (101…110).asFlow()

lifecycleScope.launch {

//TODO List 转成 Flow

mList.asFlow()

.onEach {

delay(2000)

Log.e(TAG, “onEach: $it”)

}

//TODO 组合流— 组合有点不规律

.combine(flowOther){a,b->

“本流a:其他流a:其他流a:其他流b”

}

.collect {

Log.e(TAG, “collect: $it”)

Logo

开源鸿蒙跨平台开发社区汇聚开发者与厂商,共建“一次开发,多端部署”的开源生态,致力于降低跨端开发门槛,推动万物智联创新。

更多推荐