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
}