Siddhi处理ActiveMQ无固定结构数据:类型兼容问题及条件处理咨询
解决Siddhi处理ActiveMQ无固定结构数据的类型不兼容问题
当然可以通过条件处理搞定这个问题!核心思路是先统一接收所有原始数据,再根据数据特征(比如你流里的type字段)分流并做类型转换,避免一开始就因类型不匹配报错。
具体步骤如下:
1. 定义通用原始数据流,统一接收所有数据
先把从ActiveMQ拿到的无结构数据全部用字符串类型接收,这样不管原始数据是什么类型,都能先导入Siddhi而不报错:
@source(type='jms', @map(type='csv', delimiter=',', fail.on.unknown.attribute='false', @attributes( type = '0', -- CSV第1个字段对应type time = '1', -- CSV第2个字段对应time studentId = '2', -- CSV第3个字段对应studentId extraField1 = '3',-- CSV第4个字段(对应fileId/taskId) extraField2 = '4' -- CSV第5个字段(对应totalAccesses/deadline) )), factory.initial='org.apache.activemq.jndi.ActiveMQInitialContextFactory', provider.url='tcp://127.0.0.1:61616', destination='simulatedData', connection.factory.type='queue', connection.factory.jndi.name='QueueConnectionFactory', transport.jms.SubscriptionDurable='true', transport.jms.DurableSubscriberClientID='wso2SPclient1') define stream RawSourceStream(type string, time string, studentId string, extraField1 string, extraField2 string);
这里用CSV的位置索引(0、1...)绑定字段,确保只要字段顺序对应,不管数据结构怎么变都能正确接收。
2. 根据条件分流并转换类型
利用你数据里的type字段做过滤,把不同类型的数据分流到目标流,同时用cast()函数完成类型转换:
-- 处理File类型数据,转换totalAccesses为float @info(name='route-file-data') from RawSourceStream[type == 'file'] select type, cast(time as long) as time, studentId, extraField1 as fileId, cast(extraField2 as float, 0.0) as totalAccesses -- 加默认值避免转换失败中断 insert into FileSourceProductionStream; -- 处理Task类型数据,转换deadline为long @info(name='route-task-data') from RawSourceStream[type == 'task'] select type, cast(time as long) as time, studentId, extraField1 as taskId, cast(extraField2 as long, 0) as deadline -- 转换失败用默认值0 insert into TaskSourceProductionStream; -- 定义最终的目标流(类型完全匹配) define stream FileSourceProductionStream(type string, time long, studentId string, fileId string, totalAccesses float); define stream TaskSourceProductionStream(type string, time long, studentId string, taskId string, deadline long);
关键细节说明:
- 类型转换容错:
cast()函数支持添加默认值(比如cast(xxx as float, 0.0)),如果原始数据无法转换成目标类型,会自动用默认值替代,不会导致流处理中断。 - 灵活过滤规则:如果你的数据没有
type这类明确标识,可以根据字段长度、特定值特征来过滤(比如判断第5个字段是否为数字格式)。 - 忽略未知字段:已经设置
fail.on.unknown.attribute='false',如果原始数据包含额外字段,会自动忽略,不影响正常处理。
这样处理后,就能完美解决无固定结构数据的类型不兼容问题啦!
内容的提问来源于stack exchange,提问作者mtg1989
相关产品推荐
相关产品推荐

