• 消息处理器

    将方法将消息放入队列中,消息处理器会会按照消息的入队顺序依次处理消息。

    使用场景:需要保证消息处理顺序的场景

    主要功能:

    1. 缓存接收到的消息数据
    2. 按照消息接收的顺序依次处理
    3. 支持异步消息消费
    4. 提供消息处理完成的通知机制

    Type Parameters

    • T

    Returns {
        endEnqueue: () => void;
        enqueue: (data: T) => Promise<void>;
        launch: () => void;
        onEnd: (callback: () => void) => void;
        promiseDone: () => Promise<void>;
        refistConsumer: (consumer: Consumer<T>) => void;
        termination: () => void;
    }

    返回包含以下方法的对象:

    • enqueue: 入队消息
    • refistConsumer: 注册消息消费者,接受一个promise回调函数,promise完成时才会执行下一个消息;回调函数接受一个data字段的参数,参数为入队时的数据;
    • launch: 启动消息处理
    • endEnqueue: 结束消息入队
    • onEnd: 注册处理完成回调
    • promiseDone: 返回一个promise,当所有消息处理完成时,promise状态变为resolve
    const { enqueue, refistConsumer, launch, endEnqueue, promiseDone } = messageQueue()

    // 注册消息消费者
    refistConsumer(async ({ data }) => {
    // 处理消息
    await processData(data)
    console.log('处理消息数据', data) // { step1: { content: '第一步内容' } }
    })

    // 启动消息处理
    launch()

    // 入队消息
    enqueue({ step1: { content: '第一步内容' } })
    enqueue({ step2: { content: '第二步内容' } })

    // 结束入队
    endEnqueue()

    // 等待所有消息处理完成
    await promiseDone()

    逻辑流程图

    Hero