File size: 5,901 Bytes
9ae1216 | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 | /**
* @import {BaseParquetReadOptions} from '../src/types.js'
*/
import { parquetMetadataAsync, parquetSchema } from './metadata.js'
import { parquetReadColumn, parquetReadObjects } from './read.js'
/**
* Wraps parquetRead with orderBy support.
* This is a parquet-aware query engine that can read a subset of rows and columns.
* Accepts optional orderBy column name to sort the results.
* Note that using orderBy may SIGNIFICANTLY increase the query time.
*
* @param {BaseParquetReadOptions & { orderBy?: string }} options
* @returns {Promise<Record<string, any>[]>} resolves when all requested rows and columns are parsed
*/
export async function parquetQuery(options) {
if (!options.file || !(options.file.byteLength >= 0)) {
throw new Error('parquet expected AsyncBuffer')
}
options.metadata ??= await parquetMetadataAsync(options.file, options)
const { metadata, rowStart = 0, columns, orderBy, filter } = options
if (rowStart < 0) throw new Error('parquet rowStart must be positive')
const rowEnd = options.rowEnd ?? Number(metadata.num_rows)
// Validate orderBy column exists
if (orderBy) {
const allColumns = parquetSchema(options.metadata).children.map(c => c.element.name)
if (!allColumns.includes(orderBy)) {
throw new Error(`parquet orderBy column not found: ${orderBy}`)
}
}
if (filter && !orderBy && rowEnd < metadata.num_rows) {
// iterate through row groups and filter until we have enough rows
/** @type {Record<string, any>[]} */
const filteredRows = []
let groupStart = 0
for (const group of metadata.row_groups) {
const groupEnd = groupStart + Number(group.num_rows)
// TODO: if expected > group size, start fetching next groups
const groupData = await parquetReadObjects({
...options, rowStart: groupStart, rowEnd: groupEnd,
})
filteredRows.push(...groupData)
if (filteredRows.length >= rowEnd) break
groupStart = groupEnd
}
return filteredRows.slice(rowStart, rowEnd)
} else if (filter && orderBy) {
// read all rows with orderBy column included for sorting
const readColumns = columns && !columns.includes(orderBy)
? [...columns, orderBy]
: columns
const results = await parquetReadObjects({
...options, rowStart: undefined, rowEnd: undefined, columns: readColumns,
})
// sort by orderBy column
results.sort((a, b) => compare(a[orderBy], b[orderBy]))
// project out orderBy column if not originally requested
if (readColumns !== columns) {
for (const row of results) {
delete row[orderBy]
}
}
return results.slice(rowStart, rowEnd)
} else if (filter) {
// filter without orderBy, read all matching rows
const results = await parquetReadObjects({
...options, rowStart: undefined, rowEnd: undefined,
})
return results.slice(rowStart, rowEnd)
} else if (typeof orderBy === 'string') {
// sorted but unfiltered: fetch orderBy column first
const orderColumn = await parquetReadColumn({
...options, rowStart: undefined, rowEnd: undefined, columns: [orderBy],
})
// compute row groups to fetch
const sortedIndices = Array.from(orderColumn, (_, index) => index)
.sort((a, b) => compare(orderColumn[a], orderColumn[b]))
.slice(rowStart, rowEnd)
const sparseData = await parquetReadRows({ ...options, rows: sortedIndices })
// warning: the type Record<string, any> & {__index__: number})[] is simplified into Record<string, any>[]
// when returning. The data contains the __index__ property, but it's not exposed as such.
const data = sortedIndices.map(index => sparseData[index])
return data
} else {
return await parquetReadObjects(options)
}
}
/**
* Reads a list rows from a parquet file, reading only the row groups that contain the rows.
* Returns a sparse array of rows.
* @param {BaseParquetReadOptions & { rows: number[] }} options
* @returns {Promise<(Record<string, any> & {__index__: number})[]>}
*/
async function parquetReadRows(options) {
const { file, rows } = options
options.metadata ??= await parquetMetadataAsync(file, options)
const { row_groups: rowGroups } = options.metadata
// Compute row groups to fetch
const groupIncluded = Array(rowGroups.length).fill(false)
let groupStart = 0
const groupEnds = rowGroups.map(group => groupStart += Number(group.num_rows))
for (const index of rows) {
const groupIndex = groupEnds.findIndex(end => index < end)
groupIncluded[groupIndex] = true
}
// Compute row ranges to fetch
const rowRanges = []
let rangeStart
groupStart = 0
for (let i = 0; i < groupIncluded.length; i++) {
const groupEnd = groupStart + Number(rowGroups[i].num_rows)
if (groupIncluded[i]) {
if (rangeStart === undefined) {
rangeStart = groupStart
}
} else {
if (rangeStart !== undefined) {
rowRanges.push([rangeStart, groupEnd])
rangeStart = undefined
}
}
groupStart = groupEnd
}
if (rangeStart !== undefined) {
rowRanges.push([rangeStart, groupStart])
}
// Fetch by row group and map to rows
/** @type {(Record<string, any> & {__index__: number})[]} */
const sparseData = Array(Number(options.metadata.num_rows))
for (const [rangeStart, rangeEnd] of rowRanges) {
// TODO: fetch in parallel
const groupData = await parquetReadObjects({ ...options, rowStart: rangeStart, rowEnd: rangeEnd })
for (let i = rangeStart; i < rangeEnd; i++) {
// warning: if the row contains a column named __index__, it will overwrite the index.
sparseData[i] = { __index__: i, ...groupData[i - rangeStart] }
}
}
return sparseData
}
/**
* @param {any} a
* @param {any} b
* @returns {number}
*/
function compare(a, b) {
if (a < b) return -1
if (a > b) return 1
return 0 // TODO: null handling
}
|