《數(shù)據(jù)采集與處理技術(shù)》課件-4.6 使用Python操作Kafka_第1頁
《數(shù)據(jù)采集與處理技術(shù)》課件-4.6 使用Python操作Kafka_第2頁
《數(shù)據(jù)采集與處理技術(shù)》課件-4.6 使用Python操作Kafka_第3頁
《數(shù)據(jù)采集與處理技術(shù)》課件-4.6 使用Python操作Kafka_第4頁
《數(shù)據(jù)采集與處理技術(shù)》課件-4.6 使用Python操作Kafka_第5頁
已閱讀5頁,還剩10頁未讀 繼續(xù)免費閱讀

下載本文檔

版權(quán)說明:本文檔由用戶提供并上傳,收益歸屬內(nèi)容提供方,若內(nèi)容存在侵權(quán),請進行舉報或認領(lǐng)

文檔簡介

第4章分布式消息系統(tǒng)Kafka目

錄4.1Kafka簡介4.2Kafka在大數(shù)據(jù)生態(tài)系統(tǒng)中的作用4.3Kafka與Flume的區(qū)別與聯(lián)系4.4Kafka相關(guān)概念4.5Kafka的安裝和使用4.6使用Python操作Kafka4.7Kafka與MySQL的組合使用4.6使用Python操作Kafka4.6使用Python操作Kafka使用Python操作Kafka之前,需要安裝第三方模塊python-kafka,命令如下:>pipinstallkafka-python安裝結(jié)束以后,可以使用如下命令查看已經(jīng)安裝的kafka-python的版本信息:>piplist這個命令會顯示已經(jīng)安裝的Pyhon第三方模塊,并且給出每個模塊的版本信息。4.6使用Python操作Kafka編寫一個生產(chǎn)者程序producer_test.py用來生成消息:fromkafkaimportKafkaProducer

producer=KafkaProducer(bootstrap_servers='localhost:9092')#連接Kafka

msg="HelloWorld".encode('utf-8')#發(fā)送內(nèi)容,必須是bytes類型producer.send('test',msg)#發(fā)送的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服務(wù)和Kafka服務(wù),然后,先執(zhí)行producer_test.py,再執(zhí)行consumer_test.py,就可以看到屏幕上打印出“HelloWorld”。4.6使用Python操作Kafka下面再給出一個稍微復雜一點的實例。假設(shè)有一個文件score.csv,其內(nèi)容如下:"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要求完成的任務(wù)是,Kafka生產(chǎn)者讀取文件中的所有內(nèi)容,然后,以JSON字符串的形式發(fā)送給Kafka消費者,消費者獲得消息以后轉(zhuǎn)換成表格形式打印到屏幕上,如下所示:NameScore0ZhangSan99.01LiSi45.52WangHong82.53LiuQian76.04MaLi62.55ShenTeng78.06PuWen86.54.6使用Python操作Kafka為了完成上述任務(wù),可以編寫代碼文件kafka_demo.py(要求和文件score.csv在同一個目錄下),其內(nèi)容如下:#kafka_demo.pyimportsysimportjsonimportpandasaspdimportosfromkafkaimportKafkaProducerfromkafkaimportKafkaConsumerfromkafka.errorsimportKafkaError

KAFKA_HOST="localhost"#服務(wù)器地址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":#生產(chǎn)模塊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)##字符串轉(zhuǎn)換成字典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. 本站的文檔不包含任何第三方提供的附件圖紙等,如果需要附件,請聯(lián)系上傳者。文件的所有權(quán)益歸上傳用戶所有。
  • 3. 本站RAR壓縮包中若帶圖紙,網(wǎng)頁內(nèi)容里面會有圖紙預覽,若沒有圖紙預覽就沒有圖紙。
  • 4. 未經(jīng)權(quán)益所有人同意不得將文件中的內(nèi)容挪作商業(yè)或盈利用途。
  • 5. 人人文庫網(wǎng)僅提供信息存儲空間,僅對用戶上傳內(nèi)容的表現(xiàn)方式做保護處理,對用戶上傳分享的文檔內(nèi)容本身不做任何修改或編輯,并不能對任何下載內(nèi)容負責。
  • 6. 下載文件中如有侵權(quán)或不適當內(nèi)容,請與我們聯(lián)系,我們立即糾正。
  • 7. 本站不保證下載資源的準確性、安全性和完整性, 同時也不承擔用戶因使用這些下載資源對自己和他人造成任何形式的傷害或損失。

評論

0/150

提交評論