最新国产好看的视频,伊人天堂AV在线,国产Aaaaaa视频,蜜臀视频在线观看一区,人妻av色图,密臀久久久精品影片,青青视频免费观看毛片,久草在线观看视,国产三级精品色情在线

在python環(huán)境下運用kafka對數(shù)據(jù)進行實時傳輸?shù)姆椒?/h1>
 更新時間:2018年12月27日 10:37:06   作者:真夢行路  
今天小編就為大家分享一篇在python環(huán)境下運用kafka對數(shù)據(jù)進行實時傳輸?shù)姆椒?,具有很好的參考價值,希望對大家有所幫助。一起跟隨小編過來看看吧

背景:

為了滿足各個平臺間數(shù)據(jù)的傳輸,以及能確保歷史性和實時性。先選用kafka作為不同平臺數(shù)據(jù)傳輸?shù)闹修D(zhuǎn)站,來滿足我們對跨平臺數(shù)據(jù)發(fā)送與接收的需要。

kafka簡介:

Kafka is a distributed,partitioned,replicated commit logservice。它提供了類似于JMS的特性,但是在設(shè)計實現(xiàn)上完全不同,此外它并不是JMS規(guī)范的實現(xiàn)。kafka對消息保存時根據(jù)Topic進行歸類,發(fā)送消息者成為Producer,消息接受者成為Consumer,此外kafka集群有多個kafka實例組成,每個實例(server)成為broker。無論是kafka集群,還是producer和consumer都依賴于zookeeper來保證系統(tǒng)可用性集群保存一些meta信息。

總之:kafka做為中轉(zhuǎn)站有以下功能:

1.生產(chǎn)者(產(chǎn)生數(shù)據(jù)或者說是從外部接收數(shù)據(jù))

2.消費著(將接收到的數(shù)據(jù)轉(zhuǎn)花為自己所需用的格式)

環(huán)境:

1.python3.5.x

2.kafka1.4.3

3.pandas

準(zhǔn)備開始:

1.kafka的安裝

pip install kafka-python

python環(huán)境下運用kafka對數(shù)據(jù)進行傳輸

2.檢驗kafka是否安裝成功

python環(huán)境下運用kafka對數(shù)據(jù)進行傳輸

3.pandas的安裝

pip install pandas

4.kafka數(shù)據(jù)的傳輸

直接擼代碼:

# -*- coding: utf-8 -*-
'''
@author: 真夢行路
@file: kafka.py
@time: 2018/9/3 10:20
'''
import sys
import json
import pandas as pd
import os
from kafka import KafkaProducer
from kafka import KafkaConsumer
from kafka.errors import KafkaError
 
KAFAKA_HOST = "xxx.xxx.x.xxx" #服務(wù)器端口地址
KAFAKA_PORT = 9092    #端口號
KAFAKA_TOPIC = "topic0"  #topic
 
data=pd.read_csv(os.getcwd()+'\\data\\1.csv')
key_value=data.to_json()
class Kafka_producer():
 '''
 生產(chǎn)模塊:根據(jù)不同的key,區(qū)分消息
 '''
 
 def __init__(self, kafkahost, kafkaport, kafkatopic, key):
  self.kafkaHost = kafkahost
  self.kafkaPort = kafkaport
  self.kafkatopic = kafkatopic
  self.key = key
  self.producer = KafkaProducer(bootstrap_servers='{kafka_host}:{kafka_port}'.format(
   kafka_host=self.kafkaHost,
   kafka_port=self.kafkaPort)
  )
 
 def sendjsondata(self, params):
  try:
   parmas_message = params  #注意dumps
   producer = self.producer
   producer.send(self.kafkatopic, key=self.key, value=parmas_message.encode('utf-8'))
   producer.flush()
  except KafkaError as e:
   print(e)
 
 
class Kafka_consumer():
 
 
 def __init__(self, kafkahost, kafkaport, kafkatopic, groupid,key):
  self.kafkaHost = kafkahost
  self.kafkaPort = kafkaport
  self.kafkatopic = kafkatopic
  self.groupid = groupid
  self.key = key
  self.consumer = KafkaConsumer(self.kafkatopic, group_id=self.groupid,
          bootstrap_servers='{kafka_host}:{kafka_port}'.format(
           kafka_host=self.kafkaHost,
           kafka_port=self.kafkaPort)
          )
 
 def consume_data(self):
  try:
   for message in self.consumer:
    yield message
  except KeyboardInterrupt as e:
   print(e)
 
def sortedDictValues(adict):
 items = adict.items()
 items=sorted(items,reverse=False)
 return [value for key, value in items]
 
def main(xtype, group, key):
 '''
 測試consumer和producer
 '''
 if xtype == "p":
  # 生產(chǎn)模塊
  producer = Kafka_producer(KAFAKA_HOST, KAFAKA_PORT, KAFAKA_TOPIC, key)
  print("===========> producer:", producer)
  params =key_value
  producer.sendjsondata(params)
 
 
 if xtype == 'c':
  # 消費模塊
  consumer = Kafka_consumer(KAFAKA_HOST, KAFAKA_PORT, KAFAKA_TOPIC, group,key)
  print("===========> consumer:", consumer)
 
  message = consumer.consume_data()
  for msg in message:
   msg=msg.value.decode('utf-8')
   python_data=json.loads(msg) ##這是一個字典
   key_list=list(python_data)
   test_data=pd.DataFrame()
   for index in key_list:
    print(index)
    if index=='Month':
     a1=python_data[index]
     data1 = sortedDictValues(a1)
     test_data[index]=data1
    else:
     a2 = python_data[index]
     data2 = sortedDictValues(a2)
     test_data[index] = data2
     print(test_data)
 
 
 
   # print('value---------------->', python_data)
   # print('msg---------------->', msg)
   # print('key---------------->', msg.kry)
   # print('offset---------------->', msg.offset)
 
 
 
if __name__ == '__main__':
 main(xtype='p',group='py_test',key=None)
 main(xtype='c',group='py_test',key=None)

python環(huán)境下運用kafka對數(shù)據(jù)進行傳輸

數(shù)據(jù)1.csv如下所示:

python環(huán)境下運用kafka對數(shù)據(jù)進行傳輸

幾點注意:

1、一定要有一個服務(wù)器的端口地址,不要用本機的ip或者亂寫一個ip不然程序會報錯。(我開始就是拿本機ip懟了半天,總是報錯)

2、注意數(shù)據(jù)的傳輸格式以及編碼問題(二進制傳輸),數(shù)據(jù)先轉(zhuǎn)成json數(shù)據(jù)格式傳輸,然后將json格式轉(zhuǎn)為需要格式。(不是json格式的注意dumps)

例中,dataframe->json->dataframe

3、例中dict轉(zhuǎn)dataframe,也可以用簡單方法直接轉(zhuǎn)。

eg: type(data) ==>dict,data=pd.Dateframe(data)

以上這篇在python環(huán)境下運用kafka對數(shù)據(jù)進行實時傳輸?shù)姆椒ň褪切【幏窒斫o大家的全部內(nèi)容了,希望能給大家一個參考,也希望大家多多支持腳本之家。

相關(guān)文章

  • Python天氣預(yù)報采集器實現(xiàn)代碼(網(wǎng)頁爬蟲)

    Python天氣預(yù)報采集器實現(xiàn)代碼(網(wǎng)頁爬蟲)

    這個天氣預(yù)報采集是從中國天氣網(wǎng)提取廣東省內(nèi)主要城市的天氣并回顯。本來是打算采集騰訊天氣的,但是貌似它的數(shù)據(jù)是用js寫上去還是什么的,得到的html文本中不包含數(shù)據(jù),所以就算了
    2012-10-10
  • 對Python字符串中的換行符和制表符介紹

    對Python字符串中的換行符和制表符介紹

    下面小編就為大家分享一篇對Python字符串中的換行符和制表符介紹,具有很好的參考價值,希望對大家有所幫助。一起跟隨小編過來看看吧
    2018-05-05
  • 詳解Python的Django框架中的模版繼承

    詳解Python的Django框架中的模版繼承

    這篇文章主要介紹了詳解Python的Django框架中的模版繼承,就像Python中面對對象的方法繼承道理類似,需要的朋友可以參考下
    2015-07-07
  • python 隨機打亂 圖片和對應(yīng)的標(biāo)簽方法

    python 隨機打亂 圖片和對應(yīng)的標(biāo)簽方法

    今天小編就為大家分享一篇python 隨機打亂 圖片和對應(yīng)的標(biāo)簽方法,具有很好的參考價值,希望對大家有所幫助。一起跟隨小編過來看看吧
    2018-12-12
  • Django項目單字段區(qū)間查詢的實現(xiàn)

    Django項目單字段區(qū)間查詢的實現(xiàn)

    在Django項目中會碰到一些需求就是查詢某個表中的一些字段從某日到某日的數(shù)據(jù),你可以像在SQL中那樣使用SELECT語句來查找指定字段,本文就來介紹兩種方法,感興趣的可以了解一下
    2023-10-10
  • python為tornado添加recaptcha驗證碼功能

    python為tornado添加recaptcha驗證碼功能

    tornado作為微框架,并沒有自帶驗證碼組件,recaptcha是著名的驗證碼解決方案,簡單易用,被很多公司運用來防止惡意注冊和評論。tornado添加recaptchaHA非常容易
    2014-02-02
  • Python閉包和裝飾器用法實例詳解

    Python閉包和裝飾器用法實例詳解

    這篇文章主要介紹了Python閉包和裝飾器用法,結(jié)合實例形式詳細分析了Python閉包和裝飾器的相關(guān)概念、原理、使用技巧與相關(guān)操作注意事項,需要的朋友可以參考下
    2019-05-05
  • Python中的匿名函數(shù)使用簡介

    Python中的匿名函數(shù)使用簡介

    這篇文章主要介紹了Python中的匿名函數(shù)的使用,lambda是各個現(xiàn)代編程語言中的重要功能,需要的朋友可以參考下
    2015-04-04
  • pythotn條件分支與循環(huán)詳解

    pythotn條件分支與循環(huán)詳解

    這篇文章主要介紹了Python條件分支和循環(huán)用法,結(jié)合實例形式較為詳細的分析了Python邏輯運算操作符,條件分支語句,循環(huán)語句等功能與基本用法,需要的朋友可以參考下
    2021-08-08
  • python調(diào)用機器喇叭發(fā)出蜂鳴聲(Beep)的方法

    python調(diào)用機器喇叭發(fā)出蜂鳴聲(Beep)的方法

    這篇文章主要介紹了python調(diào)用機器喇叭發(fā)出蜂鳴聲(Beep)的方法,實例分析了Python調(diào)用winsound模塊的使用技巧,需要的朋友可以參考下
    2015-03-03

最新評論

石阡县| 石河子市| 石首市| 印江| 梓潼县| 台州市| 西峡县| 凤冈县| 秦皇岛市| 西城区| 太湖县| 望谟县| 藁城市| 兰考县| 德格县| 阳泉市| 涡阳县| 绥宁县| 尼玛县| 广宁县| 广昌县| 金堂县| 南漳县| 宜川县| 怀集县| 周宁县| 普宁市| 新密市| 莆田市| 苍南县| 安溪县| 南华县| 台中县| 邵武市| 定安县| 滦平县| 奇台县| 连平县| 方正县| 兴宁市| 阿鲁科尔沁旗|