| FazBrowse GitHub Viewer | Trending | | Home |
| Tools: [Download Repo ZIP] [Original HTTPS Page] |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -5,6 +5,7 @@ import type { ChunkMessage, ClientMessage, CompleteMessage, PageMessage, Parquet | |||
| 5 | 5 | import { fromToAsyncBuffer } from './utils.js' | |
| 6 | 6 | ||
| 7 | 7 | const cache = new Map<string, Promise<AsyncBuffer>>() | |
| 8 | + const aborted = new Set<number>() | ||
| 8 | 9 | ||
| 9 | 10 | function postCompleteMessage ({ queryId, rows }: Omit<CompleteMessage, 'kind'>) { | |
| 10 | 11 | self.postMessage({ kind: 'onComplete', queryId, rows }) | |
@@ -29,30 +30,41 @@ function postParquetQueryResultMessage ({ queryId, rows }: Omit<ParquetQueryReso | |||
| 29 | 30 | } | |
| 30 | 31 | ||
| 31 | 32 | self.onmessage = async ({ data }: { data: ClientMessage }) => { | |
| 33 | + if (data.kind === 'abort') { | ||
| 34 | + aborted.add(data.queryId) | ||
| 35 | + return | ||
| 36 | + } | ||
| 32 | 37 | const { queryId, from, kind, options } = data | |
| 33 | 38 | const file = await fromToAsyncBuffer(from, cache) | |
| 34 | 39 | try { | |
| 35 | 40 | if (kind === 'parquetReadObjects') { | |
| 36 | 41 | const rows = (await parquetReadObjects({ ...options, rowFormat: 'object', file, compressors, onChunk, onPage })) as Rows | |
| 42 | + if (aborted.delete(queryId)) return | ||
| 37 | 43 | postParquetReadObjectsResultMessage({ queryId, rows }) | |
| 38 | 44 | } else if (kind === 'parquetQuery') { | |
| 39 | 45 | const rows = await parquetQuery({ ...options, file, compressors, onChunk, onPage }) | |
| 46 | + if (aborted.delete(queryId)) return | ||
| 40 | 47 | postParquetQueryResultMessage({ queryId, rows }) | |
| 41 | 48 | } else { | |
| 42 | 49 | await parquetRead({ ...options, rowFormat: 'object', file, compressors, onComplete, onChunk, onPage }) | |
| 50 | + if (aborted.delete(queryId)) return | ||
| 43 | 51 | postParquetReadResultMessage({ queryId }) | |
| 44 | 52 | } | |
| 45 | 53 | } catch (error) { | |
| 54 | + if (aborted.delete(queryId)) return | ||
| 46 | 55 | postErrorMessage({ error: error as Error, queryId }) | |
| 47 | 56 | } | |
| 48 | 57 | ||
| 49 | 58 | function onComplete(rows: Rows) { | |
| 59 | + if (aborted.has(queryId)) return | ||
| 50 | 60 | postCompleteMessage({ queryId, rows }) | |
| 51 | 61 | } | |
| 52 | 62 | function onChunk(chunk: ColumnData) { | |
| 63 | + if (aborted.has(queryId)) return | ||
| 53 | 64 | postChunkMessage({ chunk, queryId }) | |
| 54 | 65 | } | |
| 55 | 66 | function onPage(page: SubColumnData) { | |
| 67 | + if (aborted.has(queryId)) return | ||
| 56 | 68 | postPageMessage({ page, queryId }) | |
| 57 | 69 | } | |
| 58 | 70 | } | |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -68,6 +68,25 @@ function getWorker() { | |||
| 68 | 68 | return worker | |
| 69 | 69 | } | |
| 70 | 70 | ||
| 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 | + | ||
| 71 | 90 | /** | |
| 72 | 91 | * Presents almost the same interface as parquetRead, but runs in a worker. | |
| 73 | 92 | * This is useful for reading large parquet files without blocking the main thread. | |
@@ -78,11 +97,12 @@ function getWorker() { | |||
| 78 | 97 | * Note that it only supports 'rowFormat: object' (the default). | |
| 79 | 98 | */ | |
| 80 | 99 | export function parquetReadWorker(options: ParquetReadWorkerOptions): Promise<void> { | |
| 81 | - const { onComplete, onChunk, onPage, from, ...serializableOptions } = options | ||
| 100 | + const { onComplete, onChunk, onPage, from, signal, ...serializableOptions } = options | ||
| 82 | 101 | return new Promise((resolve, reject) => { | |
| 83 | 102 | const queryId = nextQueryId++ | |
| 84 | 103 | pendingAgents.set(queryId, { parquetReadResolve: resolve, reject, onComplete, onChunk, onPage }) | |
| 85 | 104 | const worker = getWorker() | |
| 105 | + if (wireAbort(worker, queryId, signal, reject)) return | ||
| 86 | 106 | const message: ClientMessage = { queryId, from, kind: 'parquetRead', options: serializableOptions } | |
| 87 | 107 | worker.postMessage(message) | |
| 88 | 108 | }) | |
@@ -98,11 +118,12 @@ export function parquetReadWorker(options: ParquetReadWorkerOptions): Promise<vo | |||
| 98 | 118 | * Note that it only supports 'rowFormat: object' (the default). | |
| 99 | 119 | */ | |
| 100 | 120 | export function parquetReadObjectsWorker(options: ParquetReadObjectsWorkerOptions): Promise<Rows> { | |
| 101 | - const { onChunk, onPage, from, ...serializableOptions } = options | ||
| 121 | + const { onChunk, onPage, from, signal, ...serializableOptions } = options | ||
| 102 | 122 | return new Promise((resolve, reject) => { | |
| 103 | 123 | const queryId = nextQueryId++ | |
| 104 | 124 | pendingAgents.set(queryId, { parquetReadObjectsResolve: resolve, reject, onChunk, onPage }) | |
| 105 | 125 | const worker = getWorker() | |
| 126 | + if (wireAbort(worker, queryId, signal, reject)) return | ||
| 106 | 127 | const message: ClientMessage = { queryId, from, kind: 'parquetReadObjects', options: serializableOptions } | |
| 107 | 128 | worker.postMessage(message) | |
| 108 | 129 | }) | |
@@ -118,11 +139,12 @@ export function parquetReadObjectsWorker(options: ParquetReadObjectsWorkerOption | |||
| 118 | 139 | * Note that it only supports 'rowFormat: object' (the default). | |
| 119 | 140 | */ | |
| 120 | 141 | export function parquetQueryWorker(options: ParquetQueryWorkerOptions): Promise<Rows> { | |
| 121 | - const { onComplete, onChunk, onPage, from, ...serializableOptions } = options | ||
| 142 | + const { onComplete, onChunk, onPage, from, signal, ...serializableOptions } = options | ||
| 122 | 143 | return new Promise((resolve, reject) => { | |
| 123 | 144 | const queryId = nextQueryId++ | |
| 124 | 145 | pendingAgents.set(queryId, { parquetQueryResolve: resolve, reject, onComplete, onChunk, onPage }) | |
| 125 | 146 | const worker = getWorker() | |
| 147 | + if (wireAbort(worker, queryId, signal, reject)) return | ||
| 126 | 148 | const message: ClientMessage = { queryId, from, kind: 'parquetQuery', options: serializableOptions } | |
| 127 | 149 | worker.postMessage(message) | |
| 128 | 150 | }) | |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -32,6 +32,12 @@ export interface ParquetReadWorkerOptions extends Omit<ParquetReadOptions, 'comp | |||
| 32 | 32 | // rowFormat 'array' is not supported in the worker. | |
| 33 | 33 | rowFormat?: 'object' | |
| 34 | 34 | 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 | ||
| 35 | 41 | } | |
| 36 | 42 | /** | |
| 37 | 43 | * Options for the worker version of parquetReadObjects | |
@@ -64,17 +70,20 @@ export interface From { | |||
| 64 | 70 | } | |
| 65 | 71 | export interface ParquetReadClientMessage extends QueryId, From { | |
| 66 | 72 | kind: 'parquetRead' | |
| 67 | - options: Omit<ParquetReadWorkerOptions, 'onComplete' | 'onChunk' | 'onPage' | 'from'> | ||
| 73 | + options: Omit<ParquetReadWorkerOptions, 'onComplete' | 'onChunk' | 'onPage' | 'from' | 'signal'> | ||
| 68 | 74 | } | |
| 69 | 75 | export interface ParquetReadObjectsClientMessage extends QueryId, From { | |
| 70 | 76 | kind: 'parquetReadObjects' | |
| 71 | - options: Omit<ParquetReadObjectsWorkerOptions, 'onChunk' | 'onPage'| 'from'> | ||
| 77 | + options: Omit<ParquetReadObjectsWorkerOptions, 'onChunk' | 'onPage'| 'from' | 'signal'> | ||
| 72 | 78 | } | |
| 73 | 79 | export interface ParquetQueryClientMessage extends QueryId, From { | |
| 74 | 80 | kind: 'parquetQuery' | |
| 75 | - options: Omit<ParquetQueryWorkerOptions, 'onComplete' | 'onChunk' | 'onPage'| 'from'> | ||
| 81 | + options: Omit<ParquetQueryWorkerOptions, 'onComplete' | 'onChunk' | 'onPage'| 'from' | 'signal'> | ||
| 76 | 82 | } | |
| 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 | ||
| 78 | 87 | ||
| 79 | 88 | /** | |
| 80 | 89 | * Messages sent by the worker to the client | |
| Back | FazBrowse Home | New Git URL |
0 commit comments