Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 5 additions & 2 deletions src/rowgroup.js
Original file line number Diff line number Diff line change
Expand Up @@ -244,8 +244,11 @@ export function assembleAsync(asyncRowGroup, schemaTree, parsers) {
/** @type {Map<string, DecodedArray>} */
const subcolumnData = new Map()
const flattened = resolved.map(({ data }) => flatten(data))
const skipped = Math.max(...resolved.map(result => result.skipped))
const end = Math.min(...resolved.map((result, i) => result.skipped + flattened[i].length))
// Physical pages can cover far more rows than the requested range.
// Clip before nested/VARIANT assembly, which expands compact binary
// values into independent object graphs and strings for every row.
const skipped = Math.max(asyncRowGroup.selectStart ?? 0, ...resolved.map(result => result.skipped))
const end = Math.min(asyncRowGroup.selectEnd ?? Infinity, ...resolved.map((result, i) => result.skipped + flattened[i].length))
for (let i = 0; i < childColumns.length; i++) {
// Offset-index reads may start each physical child at a different
// page boundary. Align them to their common absolute row range.
Expand Down
30 changes: 30 additions & 0 deletions test/rowgroup.test.js
Original file line number Diff line number Diff line change
@@ -1,9 +1,39 @@
import { describe, expect, it } from 'vitest'
import { assembleAsync, asyncGroupToRows } from '../src/rowgroup.js'
import { parquetSchema } from '../src/metadata.js'

/** @import {SchemaTree} from '../src/types.js' */

describe('assembleAsync', () => {
it('slices physical children before decoding nested variants', async () => {
const schemaTree = parquetSchema({ schema: [
{ name: 'schema', num_children: 1 },
{ name: 'payload', num_children: 2, repetition_type: 'REQUIRED', logical_type: { type: 'VARIANT' } },
{ name: 'metadata', type: 'BYTE_ARRAY', repetition_type: 'REQUIRED' },
{ name: 'value', type: 'BYTE_ARRAY', repetition_type: 'REQUIRED' },
] })
const valid = Uint8Array.from([0x11, 0x00, 0x00])
// These unselected values must never reach the VARIANT decoder. This
// also catches an implementation that decodes everything and slices later.
const invalid = Uint8Array.from([0xff])
const hi = Uint8Array.from([0x09, 0x68, 0x69])
const group = {
groupStart: 100,
groupRows: 5,
selectStart: 2,
selectEnd: 3,
asyncColumns: [
{ pathInSchema: ['payload', 'metadata'], data: Promise.resolve({ skipped: 1, data: [[invalid, valid, invalid, invalid]] }) },
{ pathInSchema: ['payload', 'value'], data: Promise.resolve({ skipped: 0, data: [[hi, hi, hi, hi, hi]] }) },
],
}
const assembled = assembleAsync(group, schemaTree)
const column = await assembled.asyncColumns[0].data
expect(column).toEqual({ skipped: 2, data: [['hi']] })
expect(await asyncGroupToRows(assembled, 2, 3, undefined, 'object'))
.toEqual([{ payload: 'hi' }])
})

it('aligns nested child columns and preserves their skipped row offset', async () => {
/** @type {SchemaTree} */
const schemaTree = {
Expand Down