Kettle实现rabbitMQ的生产与消费_rabbitmq不支持顺序消费
wptr33 2025-10-14 06:17 29 浏览
文章目录
- 一、Kettle为什么可以读取流数据?
- 二、rabbitMQ中启动MQTT插件并创建队列和路由键
- 三、Kettle实现rabbitMQ的生产与消费
Kettle是一款非常强大的ETL工具,不仅可以使用图形化界面,还可以处理各种数据,今天记录一下本人使用Kettle中MQTT组件来实现从rabbitMQ中读取流数据,并进行解析和处理。
提示:以下是本篇文章正文内容,下面案例可供参考
一、Kettle为什么可以读取流数据?
首先,本人使用的是Kettle8.2,里面关于流处理的组件有以下几种 (注意Kettle版本,我现在使用的是8.x版本,这里面只有MQTT组件,可以连接rabbitMQ,但之前使用的7.x版本是没有MQTT流处理的,也就是不能处理rabbitMQ中的数据,而9.x版本中已经有rabbitMQ组件了):
从流中获取数据信息的第一步就是第一个组件“Get records fromstream”,之后会写到,这些流处理包括JMS、Kafka、MQTT。
然后,Kettle其实是不可以直接连接rabbitMQ的,rabbitMQ默认使用amqp协议,但也可以启用MQTT插件,来使用MQTT协议。因此,我们使用Kettle通过MQTT协议步骤来生产和消费rabbitMQ。
二、rabbitMQ中启动MQTT插件并创建队列和路由键
首先使用rabbitMQ自带的控制台输入命令,也可以用windows cd到rabbit目录输入命令。
输入以下命令:
rabbitmq-plugins enable rabbitmq_mqtt 开启 rabbitmq_mqtt 对应端口 1883
rabbitmq-plugins enable rabbitmq_web_mqtt 开启 rabbitmq_web_mqtt 对应端口 15675因为我们是使用Kettle来连接rabbitMQ,所以使用的是1883端口,切记,只能使用端口1883,开启之后,可以在 http://rabbitMQ的ip地址:15672/#/ web页面查看端口是否开启:
确定端口开启之后,我们在Exchanges模块下面找到amq.topic交换器,点击进去之后,再绑定队列和路由键:
MQTT官方文档中有涉及到MQTT的系统配置,可自行尝试是否可以更改默认配置,本文未涉及:
值得注意的是我们使用的交换器只能是amq.topic,原因是rabbitMQ中的MQTT插件默认配置中只有一个交换器就是amq.topic,然后队列名称也只能是“mqtt-subscription-”开头,路由键名称可以随便设置。但要便于记忆,后续Kettle中使用的就是这个路由键。
三、Kettle实现rabbitMQ的生产与消费
1、生产数据发送给rabbitMQ
使用Kettle组件:生成记录、MQTT producer
值得注意的是,端口号只能是1883,还有就是下面的topic name是填写路由键,不是topic名称,本次绑定在amq.topic交换器下面队列的路由键是routing.update.username,所以这里填写的就是routing.update.username,其他设置默认就好,如果想要知道其他配置的作用,可参考Kettle的 官方文档 。
2、从rabbitMQ消费数据
使用Kettle组件:MQTT consumer
需要注意的还是端口和路由键,还有就是后续处理步骤最好使用英文命名,使用中文有时候会读取失败,或者识别不到XML文件,或者报错不是.ktr文件,重点切记!!!
后续处理步骤使用组件:Get records fromstream、表输出、空操作、写日志、transformation executor
“Get records fromstream”从流中接收信息,“表输出”将接收的信息存储到数据库中,“空操作”插入数据库时如果报错的消极处理,也可以换成“excel输出”,存储报错信息,“写日志”是将接收到信息打印到控制台,“transformation executor”是指定一个子转换步骤来处理数据,如后续没有处理需求,该步骤可省略,可只使用“Get records fromstream”和“写日志”两个步骤就行,进行验证。因为本次处理的数据为Json数据,所以还要对Json数据进行解析和处理,然后再使用解析后的数据去更新相关数据表。
接收到Json数据存储到了Mysql数据库中,所以解析就使用了Mysql自带的函数(JSON_EXTRACT),使用方法可参考文章:
mysql解析json字符串_Mysql解析json字符串/数组
也可参考本人的sql来解析Json数组:
select
x.id,
x.only_id,
x.createby,
x.createtime,
x.platform,
x.shopname,
x.realshopname,
x.username,
x.oldusername,
x.rownum,
y.user_id
from
(
select
a.id,
replace(json_extract(substring_index( substring_index( a.message, ";", b.id ), ";",- 1 ), '$[0].id'),'"','') as only_id,
replace(json_extract(substring_index( substring_index( a.message, ";", b.id ), ";",- 1 ), '$[0].createBy'),'"','') as createby,
from_unixtime(json_extract(substring_index( substring_index( a.message, ";", b.id ), ";",- 1 ), '$[0].createTime')/1000,'%Y-%m-%d %H:%i:%S') as createtime,
replace(json_extract(substring_index( substring_index( a.message, ";", b.id ), ";",- 1 ), '$[0].platform'),'"','') as platform,
replace(json_extract(substring_index( substring_index( a.message, ";", b.id ), ";",- 1 ), '$[0].shopName'),'"','') as shopname,
replace(json_extract(substring_index( substring_index( a.message, ";", b.id ), ";",- 1 ), '$[0].realShopName'),'"','') as realshopname,
replace(json_extract(substring_index( substring_index( a.message, ";", b.id ), ";",- 1 ), '$[0].userName'),'"','') as username,
replace(json_extract(substring_index( substring_index( a.message, ";", b.id ), ";",- 1 ), '$[0].oldUserName'),'"','') as oldusername,
b.id as rownum
from
(select id,replace(replace(replace(message,"},{","};{"),"]",""),"[","") as message,flag from sys_update_shopusername_log) a
join mysql.help_toplic_autonum b on b.id <= ( length( a.message ) - length( replace ( a.message, ";", "" ) ) + 1 )
where a.flag = 0
) x
join sys_user_detail y on x.username = y.username and upper(x.platform) = y.platform
-- “sys_update_shopusername_log”为存储的消费到的Json数据
-- “mysql.help_toplic_autonum”自定义拆分Json数组的辅助表先启动MQTT消费者,如报错,就检查ip地址、端口、路由键是否正确。启动完成后再启动MQTT生产者,发送消息给rabbitMQ,再自己消费。
消费者:
生产者:
再次查看消费者消费情况:
可以看到是能够生产数据和消费数据,这个之后就可以让上游开发将数据信息发送到我们的默认交换器amq.topic的绑定队列里面,我们就可以消费和处理了。
四、总结
注意细节:是否开启MQTT插件,端口号是否是1883,交换器和队列名称是否符合默认设定,Kettle里面MQTT producer和MQTT consumer组件所涉及到topic name 都是路由键,是在rabbitMQ中创建队列时绑定的路由键,最后就是可以根据接收到消息使用transformation executor组件来进行后续开发,转换命名最好使用英文命名。
提示:如本文有一点点帮助到您,请点赞、转发、收藏、留言,感谢!!!,如需转载、引用敬请注明!!!
相关推荐
- oracle数据导入导出_oracle数据导入导出工具
-
关于oracle的数据导入导出,这个功能的使用场景,一般是换服务环境,把原先的oracle数据导入到另外一台oracle数据库,或者导出备份使用。只不过oracle的导入导出命令不好记忆,稍稍有点复杂...
- 继续学习Python中的while true/break语句
-
上次讲到if语句的用法,大家在微信公众号问了小编很多问题,那么小编在这几种解决一下,1.else和elif是子模块,不能单独使用2.一个if语句中可以包括很多个elif语句,但结尾只能有一个else解...
- python continue和break的区别_python中break语句和continue语句的区别
-
python中循环语句经常会使用continue和break,那么这2者的区别是?continue是跳出本次循环,进行下一次循环;break是跳出整个循环;例如:...
- 简单学Python——关键字6——break和continue
-
Python退出循环,有break语句和continue语句两种实现方式。break语句和continue语句的区别:break语句作用是终止循环。continue语句作用是跳出本轮循环,继续下一次循...
- 2-1,0基础学Python之 break退出循环、 continue继续循环 多重循
-
用for循环或者while循环时,如果要在循环体内直接退出循环,可以使用break语句。比如计算1至100的整数和,我们用while来实现:sum=0x=1whileTrue...
- Python 中 break 和 continue 傻傻分不清
-
大家好啊,我是大田。今天分享一下break和continue在代码中的执行效果是什么,进一步区分出二者的区别。一、continue例1:当小明3岁时不打印年龄,其余年龄正常循环打印。可以看...
- python中的流程控制语句:continue、break 和 return使用方法
-
Python中,continue、break和return是控制流程的关键语句,用于在循环或函数中提前退出或跳过某些操作。它们的用途和区别如下:1.continue(跳过当前循环的剩余部分,进...
- L017:continue和break - 教程文案
-
continue和break在Python中,continue和break是用于控制循环(如for和while)执行流程的关键字,它们的作用如下:1.continue:跳过当前迭代,...
- 作为前端开发者,你都经历过怎样的面试?
-
已经裸辞1个月了,最近开始投简历找工作,遇到各种各样的面试,今天分享一下。其实在职的时候也做过面试官,面试官时,感觉自己问的问题很难区分候选人的能力,最好的办法就是看看候选人的github上的代码仓库...
- 面试被问 const 是否不可变?这样回答才显功底
-
作为前端开发者,我在学习ES6特性时,总被const的"善变"搞得一头雾水——为什么用const声明的数组还能push元素?为什么基本类型赋值就会报错?直到翻遍MDN文档、对着内存图反...
- 2023金九银十必看前端面试题!2w字精品!
-
导文2023金九银十必看前端面试题!金九银十黄金期来了想要跳槽的小伙伴快来看啊CSS1.请解释CSS的盒模型是什么,并描述其组成部分。答案:CSS的盒模型是用于布局和定位元素的概念。它由内容区域...
- 前端面试总结_前端面试题整理
-
记得当时大二的时候,看到实验室的学长学姐忙于各种春招,有些收获了大厂offer,有些还在苦苦面试,其实那时候的心里还蛮忐忑的,不知道自己大三的时候会是什么样的一个水平,所以从19年的寒假放完,大二下学...
- 由浅入深,66条JavaScript面试知识点(七)
-
作者:JakeZhang转发链接:https://juejin.im/post/5ef8377f6fb9a07e693a6061目录由浅入深,66条JavaScript面试知识点(一)由浅入深,66...
- 2024前端面试真题之—VUE篇_前端面试题vue2020及答案
-
添加图片注释,不超过140字(可选)1.vue的生命周期有哪些及每个生命周期做了什么?beforeCreate是newVue()之后触发的第一个钩子,在当前阶段data、methods、com...
- 今年最常见的前端面试题,你会做几道?
-
在面试或招聘前端开发人员时,期望、现实和需求之间总是存在着巨大差距。面试其实是一个交流想法的地方,挑战人们的思考方式,并客观地分析给定的问题。可以通过面试了解人们如何做出决策,了解一个人对技术和解决问...
- 一周热门
- 最近发表
- 标签列表
-
- git pull (33)
- git fetch (35)
- mysql insert (35)
- mysql distinct (37)
- concat_ws (36)
- java continue (36)
- jenkins官网 (37)
- mysql 子查询 (37)
- python元组 (33)
- mybatis 分页 (35)
- vba split (37)
- redis watch (34)
- python list sort (37)
- nvarchar2 (34)
- mysql not null (36)
- hmset (35)
- python telnet (35)
- python readlines() 方法 (36)
- munmap (35)
- docker network create (35)
- redis 集合 (37)
- python sftp (37)
- setpriority (34)
- c语言 switch (34)
- git commit (34)
