FazBrowse GitHub Viewer | Trending |
URL:
| Home
Tools: [Download Repo ZIP]   [Original HTTPS Page]

Abort parquet worker (#384) · hyparam/hyperparam-cli@98dd498 · GitHub

Commit 98dd498

Browse files
authored
Abort parquet worker (#384)
1 parent 8be8133 commit 98dd498

3 files changed

Lines changed: 50 additions & 7 deletions

File tree

‎src/lib/workers/parquetWorker.ts‎

Lines changed: 12 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -5,6 +5,7 @@ import type { ChunkMessage, ClientMessage, CompleteMessage, PageMessage, Parquet
55
import { fromToAsyncBuffer } from './utils.js'
66

77
const cache = new Map<string, Promise<AsyncBuffer>>()
8+
const aborted = new Set<number>()
89

910
function postCompleteMessage ({ queryId, rows }: Omit<CompleteMessage, 'kind'>) {
1011
self.postMessage({ kind: 'onComplete', queryId, rows })
@@ -29,30 +30,41 @@ function postParquetQueryResultMessage ({ queryId, rows }: Omit<ParquetQueryReso
2930
}
3031

3132
self.onmessage = async ({ data }: { data: ClientMessage }) => {
33+
if (data.kind === 'abort') {
34+
aborted.add(data.queryId)
35+
return
36+
}
3237
const { queryId, from, kind, options } = data
3338
const file = await fromToAsyncBuffer(from, cache)
3439
try {
3540
if (kind === 'parquetReadObjects') {
3641
const rows = (await parquetReadObjects({ ...options, rowFormat: 'object', file, compressors, onChunk, onPage })) as Rows
42+
if (aborted.delete(queryId)) return
3743
postParquetReadObjectsResultMessage({ queryId, rows })
3844
} else if (kind === 'parquetQuery') {
3945
const rows = await parquetQuery({ ...options, file, compressors, onChunk, onPage })
46+
if (aborted.delete(queryId)) return
4047
postParquetQueryResultMessage({ queryId, rows })
4148
} else {
4249
await parquetRead({ ...options, rowFormat: 'object', file, compressors, onComplete, onChunk, onPage })
50+
if (aborted.delete(queryId)) return
4351
postParquetReadResultMessage({ queryId })
4452
}
4553
} catch (error) {
54+
if (aborted.delete(queryId)) return
4655
postErrorMessage({ error: error as Error, queryId })
4756
}
4857

4958
function onComplete(rows: Rows) {
59+
if (aborted.has(queryId)) return
5060
postCompleteMessage({ queryId, rows })
5161
}
5262
function onChunk(chunk: ColumnData) {
63+
if (aborted.has(queryId)) return
5364
postChunkMessage({ chunk, queryId })
5465
}
5566
function onPage(page: SubColumnData) {
67+
if (aborted.has(queryId)) return
5668
postPageMessage({ page, queryId })
5769
}
5870
}

‎src/lib/workers/parquetWorkerClient.ts‎

Lines changed: 25 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -68,6 +68,25 @@ function getWorker() {
6868
return worker
6969
}
7070

71+
/** Wires an AbortSignal to the queryId: posts an abort message to the worker
72+
* and rejects the promise. The worker will suppress the result for this id. */
73+
function wireAbort(worker: Worker, queryId: number, signal: AbortSignal | undefined, reject: (e: Error) => void): boolean {
74+
if (!signal) return false
75+
if (signal.aborted) {
76+
pendingAgents.delete(queryId)
77+
worker.postMessage({ queryId, kind: 'abort' } satisfies ClientMessage)
78+
reject(new DOMException('Aborted', 'AbortError'))
79+
return true
80+
}
81+
signal.addEventListener('abort', () => {
82+
if (!pendingAgents.has(queryId)) return
83+
pendingAgents.delete(queryId)
84+
worker.postMessage({ queryId, kind: 'abort' } satisfies ClientMessage)
85+
reject(new DOMException('Aborted', 'AbortError'))
86+
}, { once: true })
87+
return false
88+
}
89+
7190
/**
7291
* Presents almost the same interface as parquetRead, but runs in a worker.
7392
* This is useful for reading large parquet files without blocking the main thread.
@@ -78,11 +97,12 @@ function getWorker() {
7897
* Note that it only supports 'rowFormat: object' (the default).
7998
*/
8099
export function parquetReadWorker(options: ParquetReadWorkerOptions): Promise<void> {
81-
const { onComplete, onChunk, onPage, from, ...serializableOptions } = options
100+
const { onComplete, onChunk, onPage, from, signal, ...serializableOptions } = options
82101
return new Promise((resolve, reject) => {
83102
const queryId = nextQueryId++
84103
pendingAgents.set(queryId, { parquetReadResolve: resolve, reject, onComplete, onChunk, onPage })
85104
const worker = getWorker()
105+
if (wireAbort(worker, queryId, signal, reject)) return
86106
const message: ClientMessage = { queryId, from, kind: 'parquetRead', options: serializableOptions }
87107
worker.postMessage(message)
88108
})
@@ -98,11 +118,12 @@ export function parquetReadWorker(options: ParquetReadWorkerOptions): Promise<vo
98118
* Note that it only supports 'rowFormat: object' (the default).
99119
*/
100120
export function parquetReadObjectsWorker(options: ParquetReadObjectsWorkerOptions): Promise<Rows> {
101-
const { onChunk, onPage, from, ...serializableOptions } = options
121+
const { onChunk, onPage, from, signal, ...serializableOptions } = options
102122
return new Promise((resolve, reject) => {
103123
const queryId = nextQueryId++
104124
pendingAgents.set(queryId, { parquetReadObjectsResolve: resolve, reject, onChunk, onPage })
105125
const worker = getWorker()
126+
if (wireAbort(worker, queryId, signal, reject)) return
106127
const message: ClientMessage = { queryId, from, kind: 'parquetReadObjects', options: serializableOptions }
107128
worker.postMessage(message)
108129
})
@@ -118,11 +139,12 @@ export function parquetReadObjectsWorker(options: ParquetReadObjectsWorkerOption
118139
* Note that it only supports 'rowFormat: object' (the default).
119140
*/
120141
export function parquetQueryWorker(options: ParquetQueryWorkerOptions): Promise<Rows> {
121-
const { onComplete, onChunk, onPage, from, ...serializableOptions } = options
142+
const { onComplete, onChunk, onPage, from, signal, ...serializableOptions } = options
122143
return new Promise((resolve, reject) => {
123144
const queryId = nextQueryId++
124145
pendingAgents.set(queryId, { parquetQueryResolve: resolve, reject, onComplete, onChunk, onPage })
125146
const worker = getWorker()
147+
if (wireAbort(worker, queryId, signal, reject)) return
126148
const message: ClientMessage = { queryId, from, kind: 'parquetQuery', options: serializableOptions }
127149
worker.postMessage(message)
128150
})

‎src/lib/workers/types.ts‎

Lines changed: 13 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -32,6 +32,12 @@ export interface ParquetReadWorkerOptions extends Omit<ParquetReadOptions, 'comp
3232
// rowFormat 'array' is not supported in the worker.
3333
rowFormat?: 'object'
3434
onComplete?: (rows: Rows) => void
35+
/**
36+
* Aborting the signal posts an abort message to the worker. The in-flight
37+
* read continues to completion (hyparquet has no AbortSignal support), but
38+
* the result is suppressed and the returned promise rejects with AbortError.
39+
*/
40+
signal?: AbortSignal
3541
}
3642
/**
3743
* Options for the worker version of parquetReadObjects
@@ -64,17 +70,20 @@ export interface From {
6470
}
6571
export interface ParquetReadClientMessage extends QueryId, From {
6672
kind: 'parquetRead'
67-
options: Omit<ParquetReadWorkerOptions, 'onComplete' | 'onChunk' | 'onPage' | 'from'>
73+
options: Omit<ParquetReadWorkerOptions, 'onComplete' | 'onChunk' | 'onPage' | 'from' | 'signal'>
6874
}
6975
export interface ParquetReadObjectsClientMessage extends QueryId, From {
7076
kind: 'parquetReadObjects'
71-
options: Omit<ParquetReadObjectsWorkerOptions, 'onChunk' | 'onPage'| 'from'>
77+
options: Omit<ParquetReadObjectsWorkerOptions, 'onChunk' | 'onPage'| 'from' | 'signal'>
7278
}
7379
export interface ParquetQueryClientMessage extends QueryId, From {
7480
kind: 'parquetQuery'
75-
options: Omit<ParquetQueryWorkerOptions, 'onComplete' | 'onChunk' | 'onPage'| 'from'>
81+
options: Omit<ParquetQueryWorkerOptions, 'onComplete' | 'onChunk' | 'onPage'| 'from' | 'signal'>
7682
}
77-
export type ClientMessage = ParquetQueryClientMessage | ParquetReadObjectsClientMessage | ParquetReadClientMessage
83+
export interface AbortClientMessage extends QueryId {
84+
kind: 'abort'
85+
}
86+
export type ClientMessage = ParquetQueryClientMessage | ParquetReadObjectsClientMessage | ParquetReadClientMessage | AbortClientMessage
7887

7988
/**
8089
* Messages sent by the worker to the client

0 commit comments

Comments
 (0)

Back | FazBrowse Home | New Git URL