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
8 changes: 4 additions & 4 deletions src/indexes.js
Original file line number Diff line number Diff line change
Expand Up @@ -10,17 +10,17 @@ import { deserializeTCompactProtocol } from './thrift.js'
/**
* @param {DataReader} reader
* @param {SchemaElement} schema
* @param {ParquetParsers | undefined} parsers
* @param {Partial<ParquetParsers> | undefined} parsers
* @returns {ColumnIndex}
*/
export function readColumnIndex(reader, schema, parsers = undefined) {
parsers = { ...DEFAULT_PARSERS, ...parsers }
const allParsers = { ...DEFAULT_PARSERS, ...parsers }

const thrift = deserializeTCompactProtocol(reader)
return {
null_pages: thrift.field_1,
min_values: thrift.field_2.map((/** @type {any} */ m) => convertMetadata(m, schema, parsers)),
max_values: thrift.field_3.map((/** @type {any} */ m) => convertMetadata(m, schema, parsers)),
min_values: thrift.field_2.map((/** @type {any} */ m) => convertMetadata(m, schema, allParsers)),
max_values: thrift.field_3.map((/** @type {any} */ m) => convertMetadata(m, schema, allParsers)),
boundary_order: BoundaryOrders[thrift.field_4],
null_counts: thrift.field_5,
repetition_level_histograms: thrift.field_6,
Expand Down
4 changes: 2 additions & 2 deletions src/metadata.js
Original file line number Diff line number Diff line change
Expand Up @@ -88,7 +88,7 @@ export function parquetMetadata(arrayBuffer, { parsers, geoparquet = true } = {}
const view = new DataView(arrayBuffer)

// Use default parsers if not given
parsers = { ...DEFAULT_PARSERS, ...parsers }
const allParsers = { ...DEFAULT_PARSERS, ...parsers }

// Validate footer magic number "PAR1"
if (view.byteLength < 8) {
Expand Down Expand Up @@ -148,7 +148,7 @@ export function parquetMetadata(arrayBuffer, { parsers, geoparquet = true } = {}
data_page_offset: column.field_3.field_9,
index_page_offset: column.field_3.field_10,
dictionary_page_offset: column.field_3.field_11,
statistics: convertStats(column.field_3.field_12, columnSchema[columnIndex], parsers),
statistics: convertStats(column.field_3.field_12, columnSchema[columnIndex], allParsers),
encoding_stats: column.field_3.field_13?.map((/** @type {any} */ encodingStat) => ({
page_type: PageTypes[encodingStat.field_1],
encoding: Encodings[encodingStat.field_2],
Expand Down
2 changes: 1 addition & 1 deletion src/plan.js
Original file line number Diff line number Diff line change
Expand Up @@ -295,7 +295,7 @@ export async function prefetchBloomFilters({ file, metadata, filter, filterStric
* @param {string[]} [options.columns]
* @param {Record<string, BloomFilter>[]} [options.bloomFiltersByGroup]
* @param {Record<string, SchemaElement>} [options.schemaElements]
* @param {ParquetParsers} [options.parsers]
* @param {Partial<ParquetParsers>} [options.parsers]
* @returns {Promise<{pageRangesByGroup: (PageRanges | undefined)[], pageLocationsByGroup: Record<string, PageLocation[]>[]}>}
*/
export async function prefetchPageIndexes({ file, metadata, filter, filterStrict = true, rowStart = 0, rowEnd = Infinity, columns, bloomFiltersByGroup, schemaElements, parsers }) {
Expand Down
9 changes: 5 additions & 4 deletions src/rowgroup.js
Original file line number Diff line number Diff line change
Expand Up @@ -29,9 +29,10 @@ export function readRowGroup(options, { metadata }, groupPlan) {
pathInSchema,
element: schemaPath[schemaPath.length - 1].element,
schemaPath,
parsers: { ...DEFAULT_PARSERS, ...options.parsers },
...options,
...chunk.columnMetadata,
// merge after options, so a partial parsers object keeps the defaults
parsers: { ...DEFAULT_PARSERS, ...options.parsers },
}
const { startByte, endByte } = chunk.range

Expand Down Expand Up @@ -221,12 +222,12 @@ export async function asyncGroupToRows({ asyncColumns }, selectStart, selectEnd,
*
* @param {AsyncRowGroup} asyncRowGroup
* @param {SchemaTree} schemaTree
* @param {ParquetParsers} [parsers]
* @param {Partial<ParquetParsers>} [parsers]
* @returns {AsyncRowGroup}
*/
export function assembleAsync(asyncRowGroup, schemaTree, parsers) {
const { asyncColumns } = asyncRowGroup
parsers = { ...DEFAULT_PARSERS, ...parsers }
const allParsers = { ...DEFAULT_PARSERS, ...parsers }
/** @type {AsyncColumn[]} */
const assembled = []
for (const child of schemaTree.children) {
Expand Down Expand Up @@ -256,7 +257,7 @@ export function assembleAsync(asyncRowGroup, schemaTree, parsers) {
)
}
// assemble the column
assembleNested(subcolumnData, child, parsers)
assembleNested(subcolumnData, child, allParsers)
const assembled = subcolumnData.get(child.element.name)
if (!assembled) throw new Error('parquet column data not assembled')
return { data: [assembled], skipped }
Expand Down
4 changes: 2 additions & 2 deletions src/types.d.ts
Original file line number Diff line number Diff line change
Expand Up @@ -21,7 +21,7 @@ export interface ParquetParsers {
* Parquet Metadata options for metadata parsing
*/
export interface MetadataOptions {
parsers?: ParquetParsers // custom parsers to decode advanced types
parsers?: Partial<ParquetParsers> // custom parsers to decode advanced types, merged over the defaults
geoparquet?: boolean // parse geoparquet metadata and set logical type to geometry/geography for geospatial columns (default true)
}

Expand All @@ -41,7 +41,7 @@ export interface BaseParquetReadOptions {
onPage?: (chunk: SubColumnData) => void // called when a data page is parsed. pages may contain data outside the requested range.
compressors?: Compressors // custom decompressors
utf8?: boolean // decode byte arrays as utf8 strings (default true)
parsers?: ParquetParsers // custom parsers to decode advanced types
parsers?: Partial<ParquetParsers> // custom parsers to decode advanced types, merged over the defaults
geoparquet?: boolean // parse geoparquet metadata and set logical type to geometry/geography for geospatial columns (default true)
useOffsetIndex?: boolean // use offset index to limit column chunk reads when available (default false)
useBloomFilters?: boolean // fetch bloom filters to enable row-group skipping on $eq/$in predicates (default false)
Expand Down
10 changes: 10 additions & 0 deletions test/read.test.js
Original file line number Diff line number Diff line change
Expand Up @@ -452,6 +452,16 @@ describe('parquetRead', () => {
expect(counting.bytes).toBe(14334)
})

it('keeps default parsers for types a custom parser does not override', async () => {
const file = await asyncBufferFromFile('test/files/duckdb4442.parquet')
const rows = await parquetReadObjects({
file,
parsers: { stringFromBytes: () => 'custom' },
})
expect(rows[0].call_type).toBe('custom')
expect(rows[0].call_date).toEqual(new Date('2011-10-06T22:21:49.580Z'))
})

it('filter rows with parquetRead', async () => {
const file = await asyncBufferFromFile('test/files/datapage_v2.snappy.parquet')
await parquetRead({
Expand Down