Scope
ArkTs上的Scope 和Kotlin 协程Scope 类似,用于生命周期管理
export class Scope {
private abortSignalSet = new Set<ScopeAbortSignal>()
private cancelled = false
cancel(): void {
this.abortSignalSet.forEach((abortSignal): void => abortSignal.abort())
this.abortSignalSet.clear()
this.cancelled = true
}
launch(block: (abortSignal: ScopeAbortSignal) => Promise): ScopeAbortSignal {
const abortSignal = new ScopeAbortSignal()
if (this.cancelled) {
throw new ScopeCancelError("Scope has cancelled")
}
this.abortSignalSet.add(abortSignal)
block(abortSignal)
.then((result: T) => {
abortSignal.clear()
this.abortSignalSet.delete(abortSignal)
})
.catch((e: Error) => {
abortSignal.clear()
this.abortSignalSet.delete(abortSignal)
if(e instanceof ScopeCancelError) {
KLog.debug("Scope", () => `ScopeCancelError`)
} else if (e instanceof Error){
throw e
} else {
throw new Error(JSONExt.stringifySafe(e))
}
})
return abortSignal
}
}
export class ScopeAbortSignal {
private listener: Set<()=>void> = new Set()
private _isCancelled = false
addListener(listener: ()=>void) {
this.listener.add(listener)
}
removeListener(listener: ()=>void) {
this.listener.delete(listener)
}
oneshotListener(listener: ()=>void) {
const oneshot = ()=> {
listener()
this.listener.delete(oneshot)
}
this.listener.add(oneshot)
}
invokeWhenAbort(reject:(error: Object) => void) {
this.oneshotListener(()=> {
reject(new ScopeCancelError())
})
}
abort() {
this._isCancelled = true
this.listener.forEach(listener => {
listener()
})
this.clear()
}
clear() {
this.listener.clear()
}
isCancelled(): boolean {
return this._isCancelled
}
}
用起来也比较简单:
this.scope.launch(async (abortSignal) => {
await this.flow.collect(abortSignal, async (data) => {
// logic
})
})
this.scope.cancel()的时候,可以通过abortSignal感知,执行相关的暂停工作。后面主要是配合Flow 使用。
Flow
Flow 比较重要,后续响应式变成主要依赖Flow, 同时我们要实现一些flow的常用操作符,如map,filter, combine, first 等等,尽量对标Kotlin Flow。
Flow 操作符
export abstract class Flow {
abstract collect(signal: ScopeAbortSignal, collector: FlowCollector): Promise<void>
collectIn(scope: Scope): ScopeAbortSignal {
return scope.launch(async (signal: ScopeAbortSignal) => {
await this.collect(signal, async ()=> {})
})
}
transform(transform: (signal: ScopeAbortSignal, input: T) => Promise): Flow {
return new TransformMediator(this, transform)
}
map(transform: (input: T, signal: ScopeAbortSignal) => Promise): Flow {
return this.transform(async (signal: ScopeAbortSignal, input: T) => {
return transform(input, signal)
})
}
filter(predicate: (input: T, signal: ScopeAbortSignal) => Promise<boolean>): Flow {
return new FilterMediator(this, predicate)
}
async first(signal: ScopeAbortSignal,
predicate: (input: T, signal: ScopeAbortSignal) => Promise<boolean>): Promise {
let result: T | undefined
await new CollectWhileMediator(this, async (signal: ScopeAbortSignal, input: T) => {
if (await predicate(input, signal)) {
result = input
return false
}
return true
}).collect(signal, (value: T, signal: ScopeAbortSignal) => Promise.resolve())
return result as T
}
stateIn(scope: Scope, initValue: T): StateFlow {
const flow = mutableStateFlow(initValue)
scope.launch(async (signal: ScopeAbortSignal) => {
await this.collect(signal, async (value: T, signal: ScopeAbortSignal) => {
flow.value = value
})
})
return flow
}
}
贴一下TransformMediator的实现,大家可以参考一下,其他操作符的实现也都是类似操作
class TransformMediator<INPUT, OUTPUT> extends Flow<OUTPUT> {
private up: Flow<INPUT>
private transformer: (signal: ScopeAbortSignal, input: INPUT) => Promise<OUTPUT>
constructor(up: Flow, transform: (signal: ScopeAbortSignal, input: INPUT) => Promise) {
super()
this.up = up
this.transformer = transform
}
collect(signal: ScopeAbortSignal, collector: FlowCollector<OUTPUT>): Promise<void> {
return this.up.collect(signal, async (value: INPUT, signal: ScopeAbortSignal) => {
const newValue = await this.transformer(signal, value)
await collector(newValue, signal)
})
}
}
这里着重提一下first 操作符,该操作符与将flow 中断,返回目标值。为了中断flow,需要throw error, 看一下实现:
class CollectWhileMediator extends Flow {
private up: Flow
private predicate: (signal: ScopeAbortSignal, input: T) => Promise<boolean>
constructor(up: Flow, predicate: (signal: ScopeAbortSignal, input: T) => Promise<boolean> ) {
super()
this.up = up
this.predicate = predicate
}
async collect(signal: ScopeAbortSignal, collector: FlowCollector): Promise<void> {
const ownerId = util.generateRandomUUID()
try {
await this.up.collect(signal, (value: T, signal: ScopeAbortSignal) => {
return CancellablePromiseFactory.create({
abortSignal: signal,
executor: async (resolve, reject) => {
if (!(await this.predicate(signal, value))) { // 外面返回true 直接中断
reject(new CollectCancelError(ownerId))
} else {
resolve()
}
}
})
})
} catch (e) {
if (e instanceof CollectCancelError && e.message === ownerId) {
KLog.debug('CollectWhileMediator', () => `CollectCancelError`)
// 中断后在这里返回
return
} else {
throw e instanceof Error e : new Error(`Unexpected error: ${JSONExt.stringifySafe(e)}`)
}
}
}
}
最后再提一下 combine 操作符, combine 可以同时监听多个flow,其中有一个flow 变化,就会调用transform:
export function combineArray(flows: Flow[],
transform: (data: T[], signal: ScopeAbortSignal) => Promise): Flow {
return new CombineMediator(flows, transform)
}
看一下核心实现 CombineMediator:
class CombineMediator extends Flow {
private flows: Flow[]
private transformer: (data: T[], signal: ScopeAbortSignal) => Promise
private latestValue: T[]
private hasValue: boolean[]
private mediator = mutableStateFlow<[T[], boolean[]]>([[], []])
constructor(flows: Flow[], transformer: (data: T[], signal: ScopeAbortSignal) => Promise ) {
super();
this.flows = flows
this.transformer = transformer
this.latestValue = new Array(flows.length)
this.hasValue = new Array(flows.length).fill(false)
}
async collect(signal: ScopeAbortSignal, collector: FlowCollector): Promise<void> {
const n = this.flows.length;
for (let i = 0; i < n; i++) {
const f = this.flows[i];
f.collect(signal, async (value) => {
this.latestValue[i] = value;
this.hasValue[i] = true;
this.mediator.value = [Array.from(this.latestValue), Array.from(this.hasValue)]
})
}
await this.mediator.collect(signal, async (value: [T[], boolean[]], signal: ScopeAbortSignal) => {
if (value[1].every(Boolean)) {
const newValue = await this.transformer(value[0], signal)
await collector(newValue, signal)
}
})
}
}
StateFlow
StateFlow 会保存最后一次设置的值,后面我们编写代码时主要是用这个Flow
export abstract class StateFlow extends Flow {
abstract get value(): T
}
export abstract class MutableStateFlow extends StateFlow {
abstract set value(val: T)
abstract asStateFlow(): StateFlow<T>
}
class Slot {
sequence = -1
promiseResolve: () => void
cancel: boolean = false
allocateWaitPromise(): Promise<void> {
return new Promise((resolve, reject) => {
this.promiseResolve = resolve
})
}
clear() {
this.cancel = true
this.promiseResolve.()
}
}
export class StateFlowImpl extends MutableStateFlow {
private sequence = 0
private slotRecord = new Set<Slot>()
private _value: T
constructor(initialState: T) {
super()
this._value = initialState
}
set value(val: T) {
this.updateState(val)
}
get value(): T {
return this._value
}
updateState(update: T) {
if (update === this._value) {
return
}
this.sequence++
this._value = update
this.slotRecord.forEach((slot) => {
slot.promiseResolve.()
})
}
collect(signal: ScopeAbortSignal, collector: FlowCollector): Promise<void> {
const slot = new Slot()
this.slotRecord.add(slot)
return CancellablePromiseFactory.create({
abortSignal: signal,
executor: async (resolve, reject) => {
while (true) {
if (slot.cancel) {
KLog.debug(TAG, () => `cancel`)
return
}
if (slot.sequence < this.sequence) {
slot.sequence = this.sequence
Log.debug(TAG, () => `slot.sequence=${slot.sequence} this.sequence=${this.sequence}`)
try {
await collector(this.value, signal)
} catch (e) {
if (e instanceof CollectCancelError) {
throw e
} else {
if (BuildProfile.DEBUG) {
ToastUtil.toast(`collector error ${e}n ${e.stack}`)
}
Log.error(TAG, () => `collector error ${e} n${e.stack}`)
}
}
// Log.debug(TAG, () => `collector end`)
} else {
// Log.debug(TAG, () => `collector waiting`)
await slot.allocateWaitPromise()
slot.promiseResolve = undefined
// Log.debug(TAG, () => `collector wait end`)
}
}
},
onFinished: () => {
slot.clear()
this.slotRecord.delete(slot)
},
})
}
asStateFlow(): StateFlow {
return this
}
}
export function mutableStateFlow(initialValue: T): MutableStateFlow {
return new StateFlowImpl(initialValue)
}
StateFlow 的collect 是一个while(true)的循环,内部用async/await 去控制挂起等待,
可以看到我们用到了 CancellablePromiseFactory,这个是一个Promise的简单封装,代码如下:
export class CancellablePromiseFactory {
static create(param: CancellablePromiseFactoryParam): Promise {
return new Promise(async (resolve, reject) => {
param.abortSignal.invokeWhenAbort((error) => {
param.onCancelled.()
param.onFinished.()
reject(error)
})
try {
await param.executor(
(it) => {
param.onFinished.()
resolve(it)
},
(it) => {
param.onFinished.()
reject(it)
})
} catch (e) {
param.onFinished.()
reject(e)
}
})
}
}
鸿蒙State 转 Flow
DistinctState
在鸿蒙系统中,对 State 的读写操作涉及反射机制,这会带来一定的性能损耗。因此,优化代码时应尽量减少对 State 的不必要读写操作。
DistinctState 正是为了解决这一问题而设计的,其核心目标是减少冗余的状态读写。该功能的实现逻辑相对简单,读者可直接参考相关源码。同理,项目中还提供了 DistinctListState 和 DistinctMapState 等类似工具,此处不再赘述。
@ObservedV2
export class DistinctState<T> {
// _raw 不是状态变量,对其读没有性能损耗
private _raw: T
@Trace private _state: T
get raw(): T {
return this._raw
}
get state(): T {
return this._state
}
private tag: string
constructor(initValue: T, tag: string = '') {
this._raw = initValue
this._state = initValue
this.tag = tag
}
update(newValue: T): boolean {
try {
if (equals(this.raw, newValue)) {
return false
}
} catch (e) {
Log.error('DistinctState', ()=>`equals error ${e} n ${e.stack}`)
}
this._raw = newValue
this._state = newValue
return true
}
}
DeepLink
本模块提供了一个使用极其便捷的 DeepLink 工具类,旨在简化应用内跳转逻辑。
- 首先需要添加相应的配置信息:
class DeepLinkConfig{
@DeepLink('test://zooname=elephant&animal_id={animal_id}&is_new={is_new}')
openPage(
@DeeplinkOrigin() originUrl: string,
@DeeplinkParam('animal_id') animalId: string,
@DeeplinkParam('is_new', DeepLinkParamType.Boolean) isNew: boolean
){
console.log('originUrl =>', originUrl);
console.log('animalId =>', animalId);
console.log('isNew=>', isNew);
}
}
- 调用 routeDeepLink('test://zooname=elephant&animal_id=123&is_new=true') 方法,控制台将打印出如下内容:
originUrl =>test://zooname=elephant&animal_id=123&is_new=true
animalId =>123
isNew=>true
目前,该项目已正式开源,并上传至 OpenHarmony 三方中心仓,方便开发者直接集成使用:ohpm.openharmony.cn/#/cn/detail…
关于「ArkTs 开发工具套件」的介绍至此告一段落。如果大家在后续的使用过程中遇到任何问题,欢迎在评论区留言讨论。
Android工程师的kmp(kotlin/js) for harmony开发指南 这一系列文章旨在系统性地提供一套完整的 Kotlin/Js For Harmony 解决方案。在后续的文章中,我们将继续深入介绍如何复用 ViewModel、解决序列化卡顿问题、鸿蒙开发套件的最佳实践以及架构设计思路等核心内容。敬请关注,期待与大家共同交流进步!
export class ObservedDataFlow extends StateFlowImpl {
private observer: Observer // 需要 持有Observer,否则Observer被回收,就无法感知到状态变化了
readonly tag: string
constructor(calculate: () => T, tag: string) {
super(calculate())
this.tag = tag
this.observer = new Observer(() => {
const result = calculate()
this.value = result as T
KLog.debug('ObservedDataFlow', ()=>`tag=${tag} : calculate result = ${result}`)
return result as Object
})
}
}
@ObservedV2
class Observer {
private calculate: () => Object
constructor(calculate: () => Object) {
this.calculate = calculate
}
@Computed
get value(): Object {
return this.calculate()
}
}
export function flowOf(calculate: () => T, tag: string): Flow {
return new ObservedDataFlow(calculate, tag)
}
@Local name: string = ''
const flow = flowOf(()=>this.name)
scope.launch(async (abortSignal)=>{
await this.flow
.map(async (it) => it + '-mapped')
.collect(abortSignal, async (data)=>{
print(data)
})
})
snapshotFlow{}


