版权说明:本文档由用户提供并上传,收益归属内容提供方,若内容存在侵权,请进行举报或认领
文档简介
第4章分布式消息系统Kafka目
录4.1Kafka简介4.2Kafka在大数据生态系统中的作用4.3Kafka与Flume的区别与联系4.4Kafka相关概念4.5Kafka的安装和使用4.6使用Python操作Kafka4.7Kafka与MySQL的组合使用4.6使用Python操作Kafka4.6使用Python操作Kafka使用Python操作Kafka之前,需要安装第三方模块python-kafka,命令如下:>pipinstallkafka-python安装结束以后,可以使用如下命令查看已经安装的kafka-python的版本信息:>piplist这个命令会显示已经安装的Pyhon第三方模块,并且给出每个模块的版本信息。4.6使用Python操作Kafka编写一个生产者程序producer_test.py用来生成消息:fromkafkaimportKafkaProducer
producer=KafkaProducer(bootstrap_servers='localhost:9092')#连接Kafka
msg="HelloWorld".encode('utf-8')#发送内容,必须是bytes类型producer.send('test',msg)#发送的topic为testproducer.close()4.6使用Python操作Kafka编写一个消费者程序consumer_test.py用来消费消息:fromkafkaimportKafkaConsumer
consumer=KafkaConsumer('test',bootstrap_servers=['localhost:9092'],group_id=None,auto_offset_reset='smallest')formsginconsumer:recv="%s:%d:%d:key=%svalue=%s"%(msg.topic,msg.partition,msg.offset,msg.key,msg.value)print(recv)启动Zookeeper服务和Kafka服务,然后,先执行producer_test.py,再执行consumer_test.py,就可以看到屏幕上打印出“HelloWorld”。4.6使用Python操作Kafka下面再给出一个稍微复杂一点的实例。假设有一个文件score.csv,其内容如下:"Name","Score""ZhangSan",99.0"LiSi",45.5"WangHong",82.5"LiuQian",76.0"MaLi",62.5"ShenTeng",78.0"PuWen",86.54.6使用Python操作Kafka要求完成的任务是,Kafka生产者读取文件中的所有内容,然后,以JSON字符串的形式发送给Kafka消费者,消费者获得消息以后转换成表格形式打印到屏幕上,如下所示:NameScore0ZhangSan99.01LiSi45.52WangHong82.53LiuQian76.04MaLi62.55ShenTeng78.06PuWen86.54.6使用Python操作Kafka为了完成上述任务,可以编写代码文件kafka_demo.py(要求和文件score.csv在同一个目录下),其内容如下:#kafka_demo.pyimportsysimportjsonimportpandasaspdimportosfromkafkaimportKafkaProducerfromkafkaimportKafkaConsumerfromkafka.errorsimportKafkaError
KAFKA_HOST="localhost"#服务器地址KAFKA_PORT=9092#端口号KAFKA_TOPIC="topic0"#topic
data=pd.read_csv(os.getcwd()+'\\score.csv')key_value=data.to_json()4.6使用Python操作KafkaclassKafka_producer():def__init__(self,kafkahost,kafkaport,kafkatopic,key):self.kafkaHost=kafkahostself.kafkaPort=kafkaportself.kafkatopic=kafkatopicself.key=keyducer=KafkaProducer(bootstrap_servers='{kafka_host}:{kafka_port}'.format(kafka_host=self.kafkaHost,kafka_port=self.kafkaPort))defsendjsondata(self,params):try:parmas_message=paramsproducer=ducerproducer.send(self.kafkatopic,key=self.key,value=parmas_message.encode('utf-8'))producer.flush()exceptKafkaErrorase:print(e)4.6使用Python操作KafkaclassKafka_consumer():def__init__(self,kafkahost,kafkaport,kafkatopic,groupid,key):self.kafkaHost=kafkahostself.kafkaPort=kafkaportself.kafkatopic=kafkatopicself.groupid=groupidself.key=keyself.consumer=KafkaConsumer(self.kafkatopic,group_id=self.groupid,bootstrap_servers='{kafka_host}:{kafka_port}'.format(kafka_host=self.kafkaHost,kafka_port=self.kafkaPort))defconsume_data(self):try:formessageinself.consumer:yieldmessageexceptKeyboardInterruptase:print(e)4.6使用Python操作KafkadefsortedDictValues(adict):items=adict.items()items=sorted(items,reverse=False)return[valueforkey,valueinitems]4.6使用Python操作Kafkadefmain(xtype,group,key):ifxtype=="p":#生产模块producer=Kafka_producer(KAFKA_HOST,KAFKA_PORT,KAFKA_TOPIC,key)print("===========>producer:",producer)params=key_valueproducer.sendjsondata(params)ifxtype=='c':#消费模块consumer=Kafka_consumer(KAFKA_HOST,KAFKA_PORT,KAFKA_TOPIC,group,key)print("===========>consumer:",consumer)message=consumer.consume_data()formsginmessage:msg=msg.value.decode('utf-8')python_data=json.loads(msg)##字符串转换成字典key_list=list(python_data)test_data=pd.DataFrame()forindexinkey_list:ifindex=='Name':a1=python_data[index]data1=sortedDictValues(a1)test_data[index]=data1else:a2=python_data[index]data2=sor
温馨提示
- 1. 本站所有资源如无特殊说明,都需要本地电脑安装OFFICE2007和PDF阅读器。图纸软件为CAD,CAXA,PROE,UG,SolidWorks等.压缩文件请下载最新的WinRAR软件解压。
- 2. 本站的文档不包含任何第三方提供的附件图纸等,如果需要附件,请联系上传者。文件的所有权益归上传用户所有。
- 3. 本站RAR压缩包中若带图纸,网页内容里面会有图纸预览,若没有图纸预览就没有图纸。
- 4. 未经权益所有人同意不得将文件中的内容挪作商业或盈利用途。
- 5. 人人文库网仅提供信息存储空间,仅对用户上传内容的表现方式做保护处理,对用户上传分享的文档内容本身不做任何修改或编辑,并不能对任何下载内容负责。
- 6. 下载文件中如有侵权或不适当内容,请与我们联系,我们立即纠正。
- 7. 本站不保证下载资源的准确性、安全性和完整性, 同时也不承担用户因使用这些下载资源对自己和他人造成任何形式的伤害或损失。
最新文档
- 2025安徽合肥海创城市运营管理有限公司招聘笔试历年参考题库附带答案详解
- 2025天津海泰市政绿化有限公司面向社会招聘项目经理岗位2人笔试历年参考题库附带答案详解
- 2025天津地铁9号线综合站务员招聘笔试历年参考题库附带答案详解
- 2025四川爱创科技有限公司安徽分公司招聘结构设计师等岗位3人笔试历年参考题库附带答案详解
- 2025四川岳池银泰酒店管理有限公司第四批招聘中国曲艺大酒店专业管理服务人员24人笔试历年参考题库附带答案详解
- 2025四川华丰科技股份有限公司招聘销售经理岗位测试笔试历年参考题库附带答案详解
- 2025南昌西站南昌易购便利店招聘11人笔试历年参考题库附带答案详解
- 2025北京汽车集团有限公司信息中心副主任招聘2人笔试参考题库附带答案详解
- 2025内蒙古地矿集团兴安银铅冶炼有限公司招聘40人笔试历年参考题库附带答案详解
- 2025云南民爆集团有限责任公司云南民爆相关部门缺员岗位社会招聘2人笔试参考题库附带答案详解
- 微专题:突破语病题+2026届高考语文二轮复习
- 电梯线路知识培训内容课件
- 2025转让股权合同 转让股权合同范本
- 羽毛球裁判二级考试题库及答案
- 医院安全教育与培训课件
- 锂离子电池用再生黑粉编制说明
- (正式版)DB61∕T 5033-2022 《居住建筑节能设计标准》
- 公路工程质量风险识别及控制措施
- 2025年育婴师三级试题及答案
- 2025年陕西省中考数学试题【含答案、解析】
- 民间叙事理论建构-洞察及研究
评论
0/150
提交评论