forked from hyparam/icebird
-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathavro.read.js
More file actions
211 lines (204 loc) · 7.34 KB
/
Copy pathavro.read.js
File metadata and controls
211 lines (204 loc) · 7.34 KB
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
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
import { gunzip } from 'hyparquet-compressors'
import { readZigZag, readZigZagBigInt } from './avro.metadata.js'
import { parseDecimal } from 'hyparquet/src/convert.js'
/**
* Read avro data blocks.
* Should be called after avroMetadata.
*
* @param {Object} options
* @param {DataReader} options.reader
* @param {Record<string, any>} options.metadata
* @param {Uint8Array} options.syncMarker
* @returns {Record<string, any>[]}
*/
export function avroRead({ reader, metadata, syncMarker }) {
const blocks = []
while (reader.offset < reader.view.byteLength) {
let recordCount = readZigZag(reader)
// A record count of 0 signals the end of blocks
if (recordCount === 0) break
if (recordCount < 0) {
// TODO: negative count is followed by block size for array / map
recordCount = -recordCount
}
const blockSize = readZigZag(reader)
let data = new Uint8Array(reader.view.buffer, reader.view.byteOffset + reader.offset, blockSize)
reader.offset += blockSize
// Read and verify sync marker for the block
const blockSync = new Uint8Array(reader.view.buffer, reader.view.byteOffset + reader.offset, 16)
reader.offset += 16
for (let i = 0; i < 16; i++) {
if (blockSync[i] !== syncMarker[i]) {
throw new Error('sync marker does not match')
}
}
const codec = metadata['avro.codec']
// De-compress data
if (codec === 'deflate') {
data = gunzip(data)
} else if (codec !== 'null') {
throw new Error(`unsupported codec: ${codec}`)
}
// Decode according to binary or json encoding
// Loop through metadata['avro.schema'] to parse the block
const { fields } = metadata['avro.schema']
const view = new DataView(data.buffer, data.byteOffset, data.byteLength)
const dataReader = { view, offset: 0 }
for (let i = 0; i < recordCount; i++) {
/** @type {Record<string, any>} */
const obj = {}
for (const field of fields) {
const value = readType(dataReader, field.type)
obj[field.name] = value
}
blocks.push(obj)
}
}
return blocks
}
/**
* @import {DataReader} from 'hyparquet/src/types.js'
* @import {AvroType} from '../../src/types.js'
* @param {DataReader} reader
* @param {AvroType} type
* @returns {any}
*/
function readType(reader, type) {
if (type === 'null') {
return undefined
} else if (Array.isArray(type)) {
const unionIndex = readZigZag(reader)
return readType(reader, type[unionIndex])
} else if (typeof type === 'object' && type.type === 'record') {
// read recursively
/** @type {Record<string, any>} */
const obj = {}
for (const subField of type.fields) {
obj[subField.name] = readType(reader, subField.type)
}
return obj
} else if (typeof type === 'object' && type.type === 'array') {
const arr = []
while (true) {
let count = readZigZag(reader)
if (count === 0) break
if (count < 0) {
count = -count
readZigZag(reader) // block size
}
for (let i = 0; i < count; i++) {
arr.push(readType(reader, type.items))
}
}
return arr
} else if (typeof type === 'object' && type.type === 'map') {
// Avro map: repeated blocks of (count, [string key, value]...) ending with 0.
/** @type {Record<string, any>} */
const map = {}
while (true) {
let count = readZigZag(reader)
if (count === 0) break
if (count < 0) {
count = -count
readZigZag(reader) // block size in bytes
}
for (let i = 0; i < count; i++) {
const key = readType(reader, 'string')
map[key] = readType(reader, type.values)
}
}
return map
} else if (typeof type === 'object' && type.type === 'enum') {
return type.symbols[readZigZag(reader)]
} else if (typeof type === 'object' && 'logicalType' in type && type.logicalType) {
if (type.logicalType === 'date' && type.type === 'int') {
const value = readZigZag(reader)
return new Date(value * 86400000)
} else if (type.logicalType === 'time-millis' && type.type === 'int') {
return readZigZag(reader)
} else if (type.logicalType === 'time-micros' && type.type === 'long') {
return readZigZagBigInt(reader)
} else if (type.logicalType === 'timestamp-millis' && type.type === 'long') {
const value = readZigZagBigInt(reader)
return new Date(Number(value))
} else if (type.logicalType === 'timestamp-micros' && type.type === 'long') {
const value = readZigZagBigInt(reader)
return new Date(Number(value / 1000n))
} else if (type.logicalType === 'timestamp-nanos' && type.type === 'long') {
const value = readZigZagBigInt(reader)
return new Date(Number(value / 1000000n))
} else if (type.logicalType === 'decimal' && 'precision' in type) {
const bytes = type.type === 'fixed'
? readFixed(reader, type.size)
: readType(reader, type.type)
const scale = type.scale || 0
const factor = 10 ** -scale
return parseDecimal(bytes) * factor
} else if (type.logicalType === 'uuid' && type.type === 'fixed' && type.size === 16) {
return bytesToUuid(readFixed(reader, 16))
} else {
// remaining Avro logical types (local-timestamp-*, duration, big-decimal)
// aren't used by the Iceberg spec; fall through to the underlying type.
console.warn(`unknown logical type: ${type.logicalType}`)
return type.type === 'fixed'
? readFixed(reader, type.size)
: readType(reader, type.type)
}
} else if (typeof type === 'object' && type.type === 'fixed') {
return readFixed(reader, type.size)
} else if (type === 'boolean') {
const value = reader.view.getUint8(reader.offset) === 1
reader.offset++
return value
} else if (type === 'int') {
return readZigZag(reader)
} else if (type === 'long') {
return readZigZagBigInt(reader)
} else if (type === 'float') {
const value = reader.view.getFloat32(reader.offset, true)
reader.offset += 4
return value
} else if (type === 'double') {
const value = reader.view.getFloat64(reader.offset, true)
reader.offset += 8
return value
} else if (type === 'bytes') {
const length = readZigZag(reader)
const bytes = new Uint8Array(reader.view.buffer, reader.view.byteOffset + reader.offset, length)
reader.offset += length
return bytes
} else if (type === 'string') {
const length = readZigZag(reader)
const bytes = new Uint8Array(reader.view.buffer, reader.view.byteOffset + reader.offset, length)
const text = new TextDecoder().decode(bytes)
reader.offset += length
return text
} else if (typeof type === 'object' && typeof type.type === 'string') {
// Boxed primitive/named type, e.g. { "type": "string" } or { "type": "long" }.
return readType(reader, type.type)
} else {
throw new Error(`unsupported type: ${JSON.stringify(type)}`)
}
}
/**
* @param {DataReader} reader
* @param {number} size
* @returns {Uint8Array}
*/
function readFixed(reader, size) {
const bytes = new Uint8Array(reader.view.buffer, reader.view.byteOffset + reader.offset, size)
reader.offset += size
return bytes
}
/**
* @param {Uint8Array} bytes
* @returns {string}
*/
function bytesToUuid(bytes) {
let hex = ''
for (let i = 0; i < 16; i++) {
hex += bytes[i].toString(16).padStart(2, '0')
if (i === 3 || i === 5 || i === 7 || i === 9) hex += '-'
}
return hex
}