python連接kafka加載數(shù)據(jù)的項目實踐
背景:讀取TXT文件,加載到kafka中,然后通過logstash消費(fèi)kafka中的數(shù)據(jù)加載到es中
第一步:導(dǎo)入相應(yīng)的依賴包
pip install kafka-python pip install loguru pip install msgpack
第二步:編寫連接kafka的代碼
# -*- coding: utf-8 -*-
import json
import json
import msgpack
from loguru import logger
from kafka import KafkaProducer
from kafka.errors import KafkaError
def kfk_produce_1():
"""
發(fā)送 json 格式數(shù)據(jù)
:return:
"""
producer = KafkaProducer(
//連接kafka集群的配置信息
bootstrap_servers='192.168.85.109:9092',
value_serializer=lambda v: json.dumps(v).encode('utf-8')
)
//這里是你創(chuàng)建topic和打算發(fā)送數(shù)據(jù)的地方
producer.send('python_test_topic', {'key': 'value'})
kfk_produce_1()第三步:驗證是否在kafka中創(chuàng)建topic
kafka的消費(fèi)者界面上已經(jīng)出現(xiàn)了創(chuàng)建的topic,并且數(shù)據(jù)也接收到了
注意:下面的消費(fèi)者界面的按鈕,要先運(yùn)行起來,選擇好kafka環(huán)境和topic,group以后,點(diǎn)擊那個綠色的運(yùn)行按鈕,就能實時看到發(fā)送過來的消息了,??

問題記錄:

然后在使用時,報錯提示:ImportError: cannot import name ‘KafkaConsumer’
找了一會兒最后發(fā)現(xiàn)自己創(chuàng)建的文件名叫做:kafka.py,突然意識到問題出在哪里了。
原因:
簡單說就是因為,創(chuàng)建的文件名是kafka.py,這會導(dǎo)致代碼運(yùn)行時,python解釋器查找kafka的模塊時,就找到自身kafka.py了,所以就報錯。
以后寫代碼的時候,還是要注意,切記不要用關(guān)鍵字去命名文件,避免不必要的麻煩。
到此這篇關(guān)于python連接kafka加載數(shù)據(jù)的項目實踐的文章就介紹到這了,更多相關(guān)python連接kafka加載數(shù)據(jù)內(nèi)容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!
相關(guān)文章
舉例講解Django中數(shù)據(jù)模型訪問外鍵值的方法
這篇文章主要介紹了舉例講解Django中數(shù)據(jù)模型訪問外鍵值的方法,Django是最具人氣的Python web開發(fā)框架,需要的朋友可以參考下2015-07-07
python-jwt用戶認(rèn)證食用教學(xué)的實現(xiàn)方法
這篇文章主要介紹了python-jwt用戶認(rèn)證食用教學(xué)的實現(xiàn)方法,本文給大家介紹的非常詳細(xì),對大家的學(xué)習(xí)或工作具有一定的參考借鑒價值,需要的朋友可以參考下2021-01-01
從零開始學(xué)Python第八周:詳解網(wǎng)絡(luò)編程基礎(chǔ)(socket)
本篇文章主要介紹了從零開始學(xué)Python第八周:詳解網(wǎng)絡(luò)編程基礎(chǔ)(socket) ,具有一定的參考價值,有興趣的可以了解一下。2016-12-12
OpenCV進(jìn)階之鼠標(biāo)事件的回調(diào)函數(shù)使用方法詳解
鼠標(biāo)事件的回調(diào)函數(shù)使用方法是計算機(jī)視覺領(lǐng)域的核心知識點(diǎn)之一,掌握這項技能對于提升視覺算法開發(fā)效率和應(yīng)用效果至關(guān)重要,本文深入講解了OpenCV中鼠標(biāo)事件回調(diào)函數(shù)的使用方法,感興趣的小伙伴可以了解下2026-05-05

