NodeJS消费Java Avro Kafka消息解码失败问题求助
NodeJS消费Java生成的Avro Kafka消息解码失败问题
问题详情
我尝试用NodeJS消费Java应用生成的、基于Avro Schema的Kafka消息,代码如下:
const { Kafka } = require('kafkajs'); const axios = require('axios'); const avroRegistry = require('avro-schema-registry'); const avro = require('avsc'); const { Partitioners } = require('kafkajs'); ... await consumer.subscribe({ topic, fromBeginning: true }); async function handleMessage(message) { try { const messageAsString = JSON.stringify(message); const messageObject = JSON.parse(messageAsString); const avroMessage = messageObject.message.value; console.log('messageObject:', messageObject); console.log('avroMessage:', avroMessage); const buffer = Buffer.from(avroMessage.data); console.log('buffer:', buffer); // Use Avro to decode the buffer //let type = avro.Type.forSchema(yourAvroSchema); let decoded = type.fromBuffer(buffer); } catch (error) { if (error instanceof avro.AvroError) { // Handle specific Avro decoding errors console.error('Avro decoding error:', error.message); } else { console.error('Unexpected error decoding message:', error); } } } await consumer.connect(); await consumer.run({ eachMessage: handleMessage, });
目前仅能获取到avroMessage变量,但无法完成Buffer解码,出现错误:
{"level":"ERROR","timestamp":"2024-02-25T20:04:57.433Z","logger":"kafkajs","message":"[Runner] Error when calling eachMessage","topic":"my-topic","partition":0,"offset":"98","stack":"TypeError: Right-hand side of 'instanceof' is not an object\n at Runner.handleMessage [asge]
注:该消息可正常被Java消费,Java生产时使用的是org.apache.avro.generic.GenericData.Record。
解决方案
1. 初始化Avro Type实例
代码中type变量被注释未初始化,直接调用type.fromBuffer()会导致报错。需先获取对应Avro Schema生成Type实例:
- 本地Schema场景:
// 导入本地Schema文件 const yourAvroSchema = require('./path/to/your-schema.json'); // 生成Type实例 const type = avro.Type.forSchema(yourAvroSchema); - Schema Registry场景(Java生产通常会用):
// 连接Schema Registry const registry = avroRegistry('http://your-registry-url:8081'); // 读取消息头部的magic byte和schema id(Java生产的Avro消息前5字节为1字节magic + 4字节schema id) const magicByte = buffer.readUInt8(0); const schemaId = buffer.readUInt32BE(1); // 拉取对应Schema const schema = await registry.getSchemaById(schemaId); const type = avro.Type.forSchema(schema); // 解码时跳过前5字节头部 const decoded = type.fromBuffer(buffer.slice(5));
2. 简化消息处理流程
无需将message转JSON再解析,直接读取message.value即可:
async function handleMessage({ message }) { try { // 直接获取消息原始Buffer const buffer = message.value; // 后续按步骤获取Schema并解码 } catch (error) { // 错误处理逻辑 } }
3. 修复错误判断逻辑
报错中的instanceof问题,大概率是avro.AvroError未正确识别,改用错误名称判断更可靠:
catch (error) { if (error.name === 'AvroError') { console.error('Avro解码错误:', error.message); } else { console.error('解码消息时出现意外错误:', error); } }
4. 兼容Java GenericRecord配置
Java使用GenericData.Record,解码时需配置avsc兼容Java的Avro实现:
const type = avro.Type.forSchema(schema, { wrapUnions: true, // 兼容Java的Union类型处理逻辑 logicalTypes: { // 若使用日期等逻辑类型,需对应配置转换规则 'timestamp-millis': avro.LogicalType.forSchema( { type: 'long', logicalType: 'timestamp-millis' }, (val) => new Date(val) ) } });
内容的提问来源于stack exchange,提问作者Tom
相关产品推荐
相关产品推荐

