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
相关产品推荐
相关产品推荐

