How can I queue up URLSession.DataTaskPublisher requests so that only one is made at a time?

Viewed 501

In the code below, an array of app objects is used to create an array of publishers which are merged into an array release objects.

apps.map { latestRelease(app: $0) }.merge()

Here is how latest release is done.

func latestRelease(app: App) -> AnyPublisher<Release, Error> {
    do {
        let request = try requestFactory.make(.get, "apps/\(app.owner.name)/\(app.name)/releases/latest")

        return publisherFactory.make(for: request)
            .mapError{ $0 as Error }
            .map { data, _ in data }
            .decode(type: Release.self, decoder: decoder)
            .eraseToAnyPublisher()
    } catch {
        return Fail(error: error)
            .eraseToAnyPublisher()
    }
}

The network requests are done with a factory.

struct AppCenterPublisherFactory: DataTaskPublisherFactory {
    let session: URLSession

    init(session: URLSession = .shared) {
        session.configuration.httpMaximumConnectionsPerHost = 1
        self.session = session
    }

    func make(for request: URLRequest) -> URLSession.DataTaskPublisher {
        return session.dataTaskPublisher(for: request)
    }
}

The problem is the release publishers make network requests immediately. This causes the server to return 429 Too Many Requests. How can I queue up URLSession.DataTaskPublisher requests so that only one is made at a time with a delay between each request?

2 Answers

If you use append instead of merge to combine the publishers, they will run serially instead of concurrently.

For a delay, you can prepend Empty().delay(for: cooldown, scheduler: whatever) to each release publisher except the first.

func releases<S: Scheduler>(of apps: [App], scheduler: S, cooldown: S.SchedulerTimeType.Stride)
    -> AnyPublisher<Release, Error>
{
    let singles = apps.map { latestRelease(app: $0) }
    guard let first = singles.first else { return Empty().eraseToAnyPublisher() }

    let combo: AnyPublisher<Release, Error> = singles.dropFirst()
        .reduce(first.eraseToAnyPublisher(), { combo, single in
            combo
                .append(Empty().delay(for: cooldown, scheduler: scheduler))
                .append(single)
                .eraseToAnyPublisher()
        })

    return combo
}

Test:

let testApps: [App] = [
    .init(name: "Facebook", owner: .init(name: "zuck")),
    .init(name: "Kindle", owner: .init(name: "bezos")),
    .init(name: "Crossword", owner: .init(name: "shortz")),
]

print("starting at \(Date())")
let ticket = releases(of: testApps, scheduler: DispatchQueue.global(qos: .utility), cooldown: .seconds(2))
    .sink(
        receiveCompletion: { print("got \($0) at \(Date())") },
        receiveValue: { print("got \($0) at \(Date())") })

Output:

starting at 2020-05-12 17:45:17 +0000
got Release(name: "Facebook") at 2020-05-12 17:45:17 +0000
got Release(name: "Kindle") at 2020-05-12 17:45:19 +0000
got Release(name: "Crossword") at 2020-05-12 17:45:21 +0000
got finished at 2020-05-12 17:45:21 +0000

You can deliver tasks on specified scheduler using receive(on:options:). As a sample

let queue = DispatchQueue(label: "App_Queue", qos: .default)

Then change like this

publisherFactory.make(for: request).receive(on: queue) // rest of code...

Hope you get around it. Thanks, X_X

Related