位置:首页 > Kotlin > Kotlin/Js For Harmony:ArkTs开发工具套件详解

Kotlin/Js For Harmony:ArkTs开发工具套件详解

时间:2026-08-27  |  作者:多维游侠  |  阅读:0

目录

  1. Scope
  2. Flow
  3. Flow 操作符
  4. StateFlow
  5. 鸿蒙State 转 Flow
  6. DistinctState
Flow 对应的技术说明图
Flow概括Flow的核心概念、关键要点与实践提示。
Scope 对应的技术说明图
Scope概括Scope的核心概念、关键要点与实践提示。

前言

本文深入解析Kotlin/Js For Harmony项目中的ArkTs开发工具套件,重点介绍Scope、Flow及StateFlow等核心组件。Scope机制与Kotlin协程类似,通过生命周期管理实现任务控制;Flow操作符支持数据流处理;StateFlow则用于状态同步。文章还涵盖鸿蒙State转Flow、DistinctState及DeepLink功能,为开发者提供完整的异步与状态管理方案,助力高效构建跨平台应用。

Kotlin/Js For Harmony:ArkTs开发工具套件详解 的核心流程信息图
Kotlin/Js For Harmony用简体中文信息图概括Kotlin/Js For Harmony的核心流程、关键规则与实践要点。

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 工具类,旨在简化应用内跳转逻辑。

  1. 首先需要添加相应的配置信息:
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);
    }
}
  1. 调用 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{}

免责声明:文中图文均来自网络,如有侵权请联系删除,心愿游戏发布此文仅为传递信息,不代表心愿游戏认同其观点或证实其描述。

相关文章

更多

精选合集

更多

大家都在玩

热门话题

大家都在看

更多