版权说明:本文档由用户提供并上传,收益归属内容提供方,若内容存在侵权,请进行举报或认领
文档简介
第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. 本站不保证下载资源的准确性、安全性和完整性, 同时也不承担用户因使用这些下载资源对自己和他人造成任何形式的伤害或损失。
最新文档
- 绿色古风清明主题班会
- 4.4运行与维护数据库
- 阳光体育冬季长跑活动方案4篇
- 2026化工(危险化学品)企业安全隐患排查指导手册(危险化学品仓库企业专篇)
- 麻纺厂生产进度调整办法
- 2026内蒙古鄂托克旗青少年活动中心招聘1人备考题库附参考答案详解(典型题)
- 2026中国中煤能源集团有限公司春季招聘备考题库附参考答案详解(培优b卷)
- 账务处理报税模板(商业小规模)
- 2026广东中山市绩东二社区见习生招聘备考题库附参考答案详解(a卷)
- 2026甘肃甘南州舟曲县城关镇社区卫生服务中心招聘3人备考题库含答案详解(能力提升)
- 电子厂危险性作业风险辨识及防范措施
- 煎药室培训课件讲座
- 学习通《大学生就业指导》章节测试含答案
- DB53-T913-2019-地理标志产品瑞丽柠檬-云南省
- 2025-2030中国抽动秽语综合征药物行业市场现状供需分析及投资评估规划分析研究报告
- 人工智能在改善医患沟通中的运用
- 高中生《生命教育》班会课
- T-TEA 13-2024 茶园面积 术语与定义
- 外贸项目可行性研究报告
- 农行柜面培训课件
- 《矿井通风》课件
评论
0/150
提交评论