Section 5/201 menit
5. AsyncSequence & AsyncStream
5. AsyncSequence & AsyncStream
Teori: Push vs Pull dan Backpressure
Ada dua model fundamental untuk streaming data:
Push-based (Combine, NotificationCenter, delegates):
swift
Producer ──push──▶ Consumer
"Data tersedia? Saya kirim sekarang"
Masalah: Producer bisa lebih cepat dari Consumer → Consumer kewalahan (perlu backpressure)
Pull-based (AsyncSequence):
swift
Consumer ──request──▶ Producer ──yield──▶ Consumer
"Saya siap menerima data berikutnya"
AsyncSequence menggunakan model pull: consumer meminta elemen berikutnya, dan hanya menerimanya ketika siap. Ini menciptakan backpressure natural — producer tidak bisa mengirim lebih cepat dari consumer bisa proses.
AsyncSequence — Protocol
swift
// Definisi minimal AsyncSequence:
protocol AsyncSequence {
associatedtype Element
associatedtype AsyncIterator: AsyncIteratorProtocol
func makeAsyncIterator() -> AsyncIterator
}
protocol AsyncIteratorProtocol {
associatedtype Element
mutating func next() async throws -> Element?
}
Penggunaan dengan for await:
swift
// for await adalah sintaks khusus untuk AsyncSequence
for await notification in NotificationCenter.default.notifications(named: .UIKeyboardWillShow) {
let frame = notification.userInfo?[UIKeyboardFrameEndUserInfoKey] as? CGRect
adjustForKeyboard(frame: frame)
}
// Dengan filter, map, dll — AsyncSequence punya lazy operators
let recentErrors = logStream
.filter { $0.level == .error }
.prefix(10) // hanya 10 pertama
AsyncStream — Bridge dari Callback ke Async
AsyncStream adalah cara membuat AsyncSequence dari sumber data push-based:
swift
// Contoh: wrapping NotificationCenter ke AsyncStream
func keyboardNotifications() -> AsyncStream<KeyboardInfo> {
AsyncStream { continuation in
let showObserver = NotificationCenter.default.addObserver(
forName: UIResponder.keyboardWillShowNotification,
object: nil, queue: nil
) { notification in
let info = KeyboardInfo(from: notification, isShowing: true)
continuation.yield(info) // push elemen ke stream
}
let hideObserver = NotificationCenter.default.addObserver(
forName: UIResponder.keyboardWillHideNotification,
object: nil, queue: nil
) { notification in
let info = KeyboardInfo(from: notification, isShowing: false)
continuation.yield(info)
}
// Cleanup saat stream selesai/di-cancel
continuation.onTermination = { _ in
NotificationCenter.default.removeObserver(showObserver)
NotificationCenter.default.removeObserver(hideObserver)
}
}
}
// Penggunaan:
Task { @MainActor in
for await keyboard in keyboardNotifications() {
adjustLayout(for: keyboard)
}
}
AsyncThrowingStream — Dengan Error Handling
swift
func liveUpdates(for documentID: String) -> AsyncThrowingStream<Document, Error> {
AsyncThrowingStream { continuation in
let listener = FirebaseDB.collection("docs").document(documentID)
.addSnapshotListener { snapshot, error in
if let error = error {
continuation.finish(throwing: error) // stream selesai dengan error
return
}
guard let doc = try? snapshot?.data(as: Document.self) else {
return // skip invalid snapshot
}
continuation.yield(doc) // push document update
}
continuation.onTermination = { _ in
listener.remove() // cleanup Firebase listener
}
}
}
// Penggunaan:
do {
for try await document in liveUpdates(for: "doc_123") {
updateUI(with: document)
}
} catch {
showError(error)
}
Kapan Menggunakan
AsyncSequence — gunakan untuk:
- Streaming data: WebSocket messages, server-sent events
- Real-time updates: database listeners, file watchers
- Infinite sequences: timer ticks, sensor data
- Mengganti delegation pattern yang berulang
AsyncStream — gunakan ketika:
- Menjembatani callback/delegate API ke async context
- Wrapping notification-based API
- Membuat custom event emitter
Jangan gunakan ketika:
- Hanya butuh satu nilai async (gunakan
async/awaitbiasa) - Kamu butuh backpressure yang sangat presisi (gunakan custom AsyncIterator)
- Data sudah tersedia sepenuhnya (gunakan Array biasa)
Real Use Case: WebSocket Live Chat
swift
actor WebSocketClient {
enum Message: Sendable {
case text(String)
case data(Data)
case connected
case disconnected(Error?)
}
private var webSocketTask: URLSessionWebSocketTask?
private var continuation: AsyncStream<Message>.Continuation?
var messages: AsyncStream<Message> {
AsyncStream { [weak self] continuation in
self?.continuation = continuation
continuation.onTermination = { [weak self] _ in
Task { await self?.disconnect() }
}
}
}
func connect(to url: URL) async throws {
let task = URLSession.shared.webSocketTask(with: url)
webSocketTask = task
task.resume()
continuation?.yield(.connected)
await receiveLoop()
}
private func receiveLoop() async {
while let task = webSocketTask {
do {
let message = try await task.receive()
switch message {
case .string(let text): continuation?.yield(.text(text))
case .data(let data): continuation?.yield(.data(data))
@unknown default: break
}
} catch {
continuation?.yield(.disconnected(error))
continuation?.finish()
break
}
}
}
func send(_ text: String) async throws {
try await webSocketTask?.send(.string(text))
}
func disconnect() {
webSocketTask?.cancel(with: .normalClosure, reason: nil)
webSocketTask = nil
continuation?.finish()
continuation = nil
}
}
// Penggunaan di ViewModel:
@MainActor
class ChatViewModel: ObservableObject {
@Published var messages: [ChatMessage] = []
private let client = WebSocketClient()
private var receiveTask: Task<Void, Never>?
func connect(to url: URL) async throws {
try await client.connect(to: url)
receiveTask = Task {
for await message in await client.messages {
switch message {
case .text(let text):
messages.append(ChatMessage(content: text, isMine: false))
case .disconnected:
showReconnectPrompt()
default: break
}
}
}
}
func send(text: String) async throws {
messages.append(ChatMessage(content: text, isMine: true))
try await client.send(text)
}
}
struct ChatMessage: Identifiable {
let id = UUID()
let content: String
let isMine: Bool
}