You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.29 06:37:20