如何用Node.js预构建库从Azure Avro二进制响应提取指定字段
用Node.js内置模块解析Azure Query Blob返回的Avro数据提取指定字段
当调用Azure Query Blob Contents API获取CSV数据时,返回的Avro格式是二进制结构,包含schema元数据和编码后的记录。以下是仅使用Node.js内置模块(Buffer等)解析该数据并提取hostname和serial字段的实现:
核心步骤
- 解析Avro文件头(魔法数、元数据、同步标记)
- 从元数据中提取schema,定位目标字段的位置
- 解码数据块中的记录,提取指定字段
完整代码实现
// 解码Avro的变长整数(varint) function decodeVarint(buffer, offset) { let value = 0; let shift = 0; let byte; do { byte = buffer.readUInt8(offset++); value |= (byte & 0x7F) << shift; shift += 7; } while ((byte & 0x80) !== 0); return { value, offset }; } // 解析Avro二进制数据 function parseAvroData(buffer) { let offset = 0; // 验证Avro魔法数 const magic = buffer.slice(offset, offset + 4); if (!magic.equals(Buffer.from([0x4F, 0x62, 0x6A, 0x01]))) { throw new Error("无效的Avro格式数据"); } offset += 4; // 读取元数据长度并解析元数据JSON const { value: metadataLen, offset: metaOffset } = decodeVarint(buffer, offset); offset = metaOffset; const metadata = JSON.parse(buffer.slice(offset, offset + metadataLen).toString('utf8')); offset += metadataLen; // 解析schema并定位目标字段 const schema = JSON.parse(metadata.schema); const targetFields = new Set(['hostname', 'serial']); const fieldMap = schema.fields.reduce((map, field, idx) => { if (targetFields.has(field.name)) { map[field.name] = idx; } return map; }, {}); // 跳过同步标记(16字节) offset += 16; const extractedRecords = []; // 遍历所有数据块 while (offset < buffer.length) { // 读取当前块的记录数和数据长度 const { value: recordCount, offset: countOffset } = decodeVarint(buffer, offset); offset = countOffset; const { value: dataLen, offset: dataOffset } = decodeVarint(buffer, offset); offset = dataOffset; // 处理当前块的数据 const dataBlock = buffer.slice(offset, offset + dataLen); offset += dataLen; // 跳过块末尾的同步标记 offset += 16; let blockOffset = 0; for (let i = 0; i < recordCount; i++) { const record = {}; // 遍历schema所有字段,仅提取目标字段 for (const [fieldName] of Object.entries(fieldMap)) { // 读取字符串类型字段的长度和内容(适配CSV导出的字符串字段) const { value: strLen, offset: strOffset } = decodeVarint(dataBlock, blockOffset); blockOffset = strOffset; record[fieldName] = dataBlock.slice(blockOffset, blockOffset + strLen).toString('utf8'); blockOffset += strLen; } // 跳过非目标字段的内容 for (const field of schema.fields) { if (!targetFields.has(field.name)) { const { value: skipLen, offset: skipOffset } = decodeVarint(dataBlock, blockOffset); blockOffset = skipOffset; blockOffset += skipLen; } } extractedRecords.push(record); } } return extractedRecords; } // 使用示例(假设已获取到Avro响应Buffer) async function processAzureBlobResponse() { // 替换为你的响应Buffer获取逻辑,比如通过http/https模块请求后拼接的Buffer const avroBuffer = await getAzureQueryBlobResponse(); try { const results = parseAvroData(avroBuffer); results.forEach(item => { console.log(`hostname: ${item.hostname}, serial: ${item.serial}`); }); } catch (err) { console.error("解析失败:", err.message); } } // 模拟获取响应的函数(实际替换为你的API调用逻辑) async function getAzureQueryBlobResponse() { const https = require('https'); return new Promise((resolve, reject) => { const req = https.get('YOUR_AZURE_BLOB_QUERY_URL', { headers: { 'Authorization': 'YOUR_AUTH_TOKEN', // 其他必要的请求头 } }, (res) => { const chunks = []; res.on('data', (chunk) => chunks.push(chunk)); res.on('end', () => resolve(Buffer.concat(chunks))); }); req.on('error', reject); }); } // 执行处理 processAzureBlobResponse();
注意事项
- 上述代码假设
hostname和serial字段为字符串类型(CSV导出的默认类型),如果字段是整数/其他类型,需要修改对应字段的解码逻辑(比如整数直接使用decodeVarint的返回值)。 - 确保获取的响应Buffer是完整的Avro文件,没有截断,否则会导致解析失败。
- 若需要严格校验数据完整性,可以添加同步标记的对比逻辑(代码中已跳过该步骤以简化实现)。
内容的提问来源于stack exchange,提问作者Vinoth
相关产品推荐
相关产品推荐

