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

Siddhi Avro Source连接器无法解码Kafka Avro消息问题求助

问题根因

报错提示Expected byte Array or ByteBuffer, but found java.lang.String,本质是Siddhi的Kafka输入源默认会把读取到的消息载荷转换为字符串格式传递给后续的Avro映射器,但Avro反序列化要求输入必须是原始二进制字节数组,类型不匹配直接触发了映射失败异常。

解决方案

1. 调整Kafka源配置

删除原配置里的use.avro.deserializer="true",同时在Kafka源参数中新增is.binary.message="true",强制Kafka源将消息以原始字节数组的形式传递给Avro映射器,不做字符串转码处理。

2. 修正后的完整Siddhi应用代码

@App:name('UMBAlarm')

@sink(type='log')
define stream logStream(OCName string );

@source(type='kafka',
        topic.list='TEST2',
        partition.no.list='0',
        threading.option='single.thread',
        group.id="group",
        bootstrap.servers='bt1svpff:9092',
        is.binary.message="true",
        @map(type='avro',
        schema.def = """{
            "type":"record",
            "name":"AvroTemipAlarm",
            "namespace":"com.hp.ossa.fault.avro",
            "fields":[
                {"name":"OCName","type":"string"},
                {"name":"Identifier","type":"long"},
                {"name":"AttributeList","type":
                    {"type":"array","items":
                        {"type":"record",
                            "name":"AttributeRecord",
                            "fields":[
                                {"name":"AttributeId","type":"long"},
                                {"name":"AttributeName","type":"string"},
                                {"name":"AttributeType","type":"int"},
                                {"name":"IntValue","type":{"type":"array","items":"int"}},
                                {"name":"LongValue","type":{"type":"array","items":"long"}},
                                {"name":"StringValue","type":{"type":"array","items":"string"}},
                                {"name":"BooleanValue","type":{"type":"array","items":"boolean"}},
                                {"name":"DoubleValue","type":{"type":"array","items":"double"}}
                            ]
                        }
                    }
                }
            ]
        } """,
        @attributes(OCName="OCName")
    )
)
define stream hbStream (OCName string);
from hbStream select * insert into logStream;

你之前写的Python解码器可以正常运行,是因为代码直接读取Kafka消息的原始字节载荷做反序列化,没有额外的字符串转义步骤,和修改后Siddhi的处理逻辑完全一致。

内容的提问来源于stack exchange,提问作者SrikanthR

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.26 05:15:06