【问题标题】:Synchronously queueing Asynchronous operations同步排队 异步操作
【发布时间】:2019-02-07 19:09:09
【问题描述】:

我有一个显示在 UITableView 上的项目列表/数组。每个项目都包含一个状态机,它使用期货和承诺机制执行一些异步任务,并且它的状态显示在相应的 UITableViewCell 上。目前对于一个模型来说它工作得很好。

但是,我需要将其作为批处理序列进行。例如,一个数组中可以有 15 个模型,但是在特定时刻我只能启动 3 个模型,一旦模型完成或失败,我应该手动触发第 4 个模型来启动它的任务。注意:我无法启动所有 15 个模型的操作,而只是等待回调,因为它受到硬件的限制,并且在那种情况下会立即失败。

如果上面不清楚,具体来说,下面是两个例子: 我的问题陈述与 iPhone App Store 应用中 Updates 选项卡下的 全部更新 功能完全匹配。如果您有 20 个应用程序更新并点击“全部更新”按钮,它会显示 17 个处于等待状态的应用程序,并且随时仅在 3 个应用程序上运行更新下载。应用程序更新完成后,将移至下一个。这是我的问题陈述的精确复制品,但有一点小改动。

Twist:我的操作是通过蓝牙进行的硬件相关操作。想想你有 20 个可穿戴设备,你想通过蓝牙写入一些数据来进行配置。硬件限制是,您一次最多可以连接 3-4 个设备。因此,一旦设备/外围设备操作成功或失败,我应该尝试连接 4 个,依此类推,直到所有操作都完成。还有重试功能,可以将失败的队列返回。

我的问题是我应该如何构建来维护和监控。我对并发有一个大致的了解,但是并没有做太多的工作。我目前的感觉是使用包裹在 Manager 类中的队列和计数器来监视状态。想要一些关于如何解决这个问题的帮助。我也不需要代码,只需要数据结构的概念解决方案。

【问题讨论】:

  • 将每个 Job 放入一个数组中。每个 Job 都是一个闭包,每个 Job 都应该有一个完成方法。完成工作后,看看剩下的工作,然后开始下一个工作。您应该使用同步的作业数组,并且必须处理失败和取消。
  • @StephanJanuar 是的,工作、取消、完成和失败处理都在那里。我不确定把它放在一个数组中会起作用。因为它是异步的,所以可以按任何顺序完成,我相信状态管理会很麻烦。
  • 基于basememara.com/creating-thread-safe-arrays-in-swift我已经管理了所有的等待任务。但你是对的,这很麻烦。

标签: ios swift data-structures concurrency


【解决方案1】:

我认为你的情况是一个很好的例子。我用 ReactiveKit 整理了一个简单的例子,你可以看看。反应套件非常简单,对于这种情况来说绰绰有余。您也可以使用任何其他反应式库。我希望它有所帮助。

ReactiveKit:https://github.com/DeclarativeHub/ReactiveKit 邦德:https://github.com/DeclarativeHub/Bond

安装 reactiveKit 依赖后,您可以在工作区中运行以下代码:

import UIKit
import Bond
import ReactiveKit

class ViewController: UIViewController {

    var jobHandler : JobHandler!
    var jobs = [Job(name: "One", state: nil), Job(name: "Two", state: nil), Job(name: "Three", state: nil), Job(name: "Four", state: nil), Job(name: "Five", state: nil)]

    override func viewDidLoad() {
        super.viewDidLoad()
        self.jobHandler = JobHandler()
        self.run()
    }

    func run() {
        // Initialize jobs with queue state
        _ = self.jobs.map({$0.state.value = .queue})
        self.jobHandler.jobs.insert(contentsOf: jobs, at: 0)
        self.jobHandler.queueJobs(limit: 2) // Limit of how many jobs you can start with
    }
}

// Job state, I added a few states just as test cases, change as required
public enum State {
    case queue, running, completed, fail
}

class Job : Equatable {

    // Initialize state as a Reactive property
    var state = Property<State?>(nil)
    var name : String!

    init(name: String, state: State?) {
        self.state.value = state
        self.name = name
    }

    // This runs the current job
    typealias jobCompletion = (State) -> Void
    func runJob (completion: @escaping jobCompletion) {
        DispatchQueue.main.asyncAfter(deadline: .now() + 2.0) {
            self.state.value = .completed
            completion(self.state.value ?? .fail)
            return
        }
        self.state.value = .running
        completion(.running)
    }

    // To find the index of current job
    static func == (lhs: Job, rhs: Job) -> Bool {
        return lhs.name == rhs.name
    }
}

class JobHandler {

    // The array of jobs in an observable form, so you can see event on the collection
    var jobs = MutableObservableArray<Job>([])
    // Completed jobs, you can add failed jobs as well so you can queue them again
    var completedJobs = [Job]()


    func queueJobs (limit: Int) {
        // Observe the events in the datasource
        _ = self.jobs.observeNext { (collection) in
            let jobsToRun = collection.collection.filter({$0.state.value == .queue})
            self.startJob(jobs: Array(jobsToRun.prefix(limit)))
        }.dispose()
    }

    func startJob (jobs: [Job?]) {
        // Starts a job thrown by the datasource event
        jobs.forEach { (job) in
            guard let job = job else { return }
            job.runJob { (state) in
                switch state {
                case .completed:
                    if !self.jobs.collection.isEmpty {
                        guard let index = self.jobs.collection.indexes(ofItemsEqualTo: job).first else { return }
                        print("Completed " + job.name)
                        self.jobs.remove(at: index)
                        self.completedJobs.append(job)
                        self.queueJobs(limit: 1)
                    }
                case .queue:
                    print("Queue")
                case .running:
                    print("Running " + job.name)
                case .fail:
                    print("Fail")
                }
            }
        }
    }
}

extension Array where Element: Equatable {
    func indexes(ofItemsEqualTo item: Element) -> [Int]  {
        return enumerated().compactMap { $0.element == item ? $0.offset : nil }
    }
} 

【讨论】:

  • 我的代码看起来接近这个。是的,它是反应式的,但使用期货和承诺。我会用同样的方法尝试这种方法,看看效果如何。感谢您提供示例代码。
  • 嘿 Galo,感谢您的代码示例。我试图通过操作和操作队列来实现相同的目标。 :)
【解决方案2】:

使用OperationQueue,对于每种操作,您都有Operation 的子类。通过操作,您可以添加对其他操作的依赖以完成。因此,例如,如果您想等到第 3 次操作完成后再启动第 4、第 5 和第 6 次操作,您只需将第 3 次操作添加为它们的依赖项。

编辑:因此,要将操作组合在一起,您可以为其创建一个单独的类。我在下面添加了一个代码示例。 add(dependency: OperationGroup) 函数告诉其他组在原始组中执行完操作后立即开始执行操作。

//Make a subclass for each kind of operation.
class BluetoothOperation: Operation
{
    let number: Int

    init(number: Int)
    {
        self.number = number
    }

    override func main() 
    {
        print("Executed bluetooth operation number \(number)")
    }
}

class OperationGroup
{
    var operationCounter: Int = 0
    var operations: [Operation]
    let operationQueue: OperationQueue = OperationQueue()

    init(operations: [Operation])
    {
        self.operations = operations
    }

    func executeAllOperations()
    {
        operationQueue.addOperations(operations, waitUntilFinished: true)
    }

    //This in essence is popping the "Stack" of operations you have.
    func pop() -> Operation?
    {
        guard operationCounter < operations.count else { return nil }

        let operation = operations[operationCounter]

        operationCounter += 1

        return operation
    }

    func add(dependency: OperationGroup)
    {
        dependency.operations.forEach(
        {
            $0.completionBlock =
            {
                if let op = self.pop()
                {
                    dependency.operationQueue.addOperation(op)
                }
            }
        })
    }
}


let firstOperationGroup = OperationGroup(operations: [BluetoothOperation(number: 1), BluetoothOperation(number: 2), BluetoothOperation(number: 3)])
let secondOperationGroup = OperationGroup(operations: [BluetoothOperation(number: 4), BluetoothOperation(number: 5), BluetoothOperation(number: 6)])

secondOperationGroup.add(dependency: firstOperationGroup)

firstOperationGroup.executeAllOperations()

【讨论】:

  • 然而,我想到了这一点,即使它是同步的,顺序也可能是随机的。例如,3 次操作正在进行,第 2 次完成,我应该启动第 5 次操作,第 1 次失败,应该启动第 6 次 n 等等。但是,在上述方式中,它是紧密耦合的。
  • 我的意思是虽然它是同步的,但顺序是随机的,运行时间由事件驱动。
  • 当然。也感谢您的示例。将尝试。
  • 嗨 Atharva,感谢您的帮助。我尝试了一种不同的方法,使其基于 OperationQueue 动态化,并通过maxConcurrentOperationCount 处理它,并使用每个操作完成块来排队下一个(如果有的话)。 :)
  • @RameswarPrasad 使用 OperationQueue 您最好定义一个异步操作。然后只需定义最大并发操作并让它们运行。不过,操作的正确子类是 PITA。有更好的方法。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2016-08-24
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2019-05-27
  • 1970-01-01
相关资源
最近更新 更多