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

Golang監(jiān)聽日志文件并發(fā)送到kafka中

 更新時間:2022年04月14日 15:43:28   作者:zhijie  
這篇文章主要介紹了Golang監(jiān)聽日志文件并發(fā)送到kafka中,日志收集項目的準備中,本文主要講的是利用golang的tail庫,監(jiān)聽日志文件的變動,將日志信息發(fā)送到kafka中?,需要的朋友可以參考一下

前言

日志收集項目的準備中,本文主要講的是利用golang的tail庫,監(jiān)聽日志文件的變動,將日志信息發(fā)送到kafka中。

涉及的golang庫和可視化工具:

go-ini,sarama,tail其中:

  • go-ini:用于讀取配置文件,統(tǒng)一管理配置項,有利于后其的維護
  • sarama:是一個go操作kafka的客戶端。目前我用于向kefka發(fā)送消息
  • tail:類似于linux的tail命令了,讀取文件的后幾行。如果文件有追加數據,會檢測到。就是通過它來監(jiān)聽日志文件

可視化工具:

offsetexplorer:是kafka的可視化工具,這里用來查看消息是否投遞成功

工作的流程

  • 加載配置,初始化saramakafka。
  • 起一個的協(xié)程,利用tail不斷去監(jiān)聽日志文件的變化。
  • 主協(xié)程中一直阻塞等待tail發(fā)送消息,兩者通過一個管道通訊。一旦主協(xié)程接收到新日志,組裝格式,然后發(fā)送到kafka中

main.png

環(huán)境準備

環(huán)境的話,確保zookeeperkafka正常運行。因為還沒有使用sarama讀取數據,使用offsetexplorer來查看任務是否真的投遞成功了。

代碼分層

serve來存放寫tail服務類和sarama服務類,conf存放ini配置文件

main函數為程序入口

 

pro-dir.png

關鍵的代碼

main.go

main函數做的有:構建配置結構體,映射配置文件。調用和初始化tail,srama服務。

package main

import (
	"fmt"
	"sarama/serve"

	"github.com/go-ini/ini"
)

type KafkaConfig struct {
	Address     string `ini:"address"`
	ChannelSize int    `ini:"chan_size"`
}
type TailConfig struct {
	Path     string `ini:"path"`
	Filename string `ini:"fileName"`
	// 如果是結構體,則指明分區(qū)名
	Children `ini:"tailfile.children"`
}
type Config struct {
	KafkaConfig `ini:"kafka"`
	TailConfig  `ini:"tailfile"`
}
type Children struct {
	Name string `ini:"name"`
}

func main() {
	// 加載配置
	var cfg = new(Config)
	err := ini.MapTo(cfg, "./conf/go-conf.ini")
	if err != nil {
		fmt.Print(err)
	}
	// 初始化kafka
	ks := &serve.KafukaServe{}
	// 啟動kafka消息監(jiān)聽。異步
	ks.InitKafka([]string{cfg.KafkaConfig.Address}, int64(cfg.KafkaConfig.ChannelSize))
	// 關閉主協(xié)程時,關閉channel
	defer ks.Destruct()

	// 初始化tail
	ts := &serve.TailServe{}
	ts.TailInit(cfg.TailConfig.Path + "/" + cfg.TailConfig.Filename)
	// 阻塞
	ts.Listener(ks.MsgChan)

}

kafka.go

有3個方法 :

  • InitKafka,組裝配置項以及初始化接收消息的管道,
  • Listener,監(jiān)聽管道消息,收到消息后,將消息組裝,發(fā)送到kafka
  • Destruct, 關閉管道
package serve

import (
	"fmt"

	"github.com/Shopify/sarama"
)

type KafukaServe struct {
	MsgChan chan string
	//err         error
}

func (ks *KafukaServe) InitKafka(addr []string, chanSize int64) {

	// 讀取配置
	config := sarama.NewConfig()
	// 1. 初始化生產者配置
	config.Producer.RequiredAcks = sarama.WaitForAll
	// 選擇分區(qū)
	config.Producer.Partitioner = sarama.NewRandomPartitioner
	// 成功交付的信息
	config.Producer.Return.Successes = true

	ks.MsgChan = make(chan string, chanSize)

	go ks.Listener(addr, chanSize, config)

}

func (ks *KafukaServe) Listener(addr []string, chanSize int64, config *sarama.Config) {
	//  連接kafka
	var kafkaClient, _ = sarama.NewSyncProducer(addr, config)
	defer kafkaClient.Close()
	for {
		select {
		case content := <-ks.MsgChan:
			//
			msg := &sarama.ProducerMessage{
				Topic: "weblog",
				Value: sarama.StringEncoder(content),
			}
			partition, offset, err := kafkaClient.SendMessage(msg)
			if err != nil {
				fmt.Println(err)
			}
			fmt.Println("分區(qū),偏移量:")
			fmt.Println(partition, offset)
			fmt.Println("___")
		}

	}
}

func (ks *KafukaServe) Destruct() {
	close(ks.MsgChan)
}

tail.go

主要包括了兩個方法:

  • TailInit初始化,組裝tail配置。Listener
  • Listener,保存kafka服務類初始化之后的管道。監(jiān)聽日志文件,如果有新日志,就往管道里發(fā)送
package serve

import (
	"fmt"

	"github.com/hpcloud/tail"
)

type TailServe struct {
	tails *tail.Tail
}

func (ts *TailServe) TailInit(filenName string) {
	config := tail.Config{
		ReOpen:    true,
		Follow:    true,
		Location:  &tail.SeekInfo{Offset: 0, Whence: 2},
		MustExist: false,
		Poll:      true,
	}
	// 打開文件開始讀取數據

	ts.tails, _ = tail.TailFile(filenName, config)

	// if err != nil {
	// 	fmt.Println("tails %s failed,err:%v\n", filenName, err)
	// 	return nil, err
	// }
	fmt.Println("啟動," + filenName + "監(jiān)聽")
}

func (ts *TailServe) Listener(MsgChan chan string) {
	for {
		msg, ok := <-ts.tails.Lines
		if !ok {
			// todo
			fmt.Println("數據接收失敗")
			return
		}
		fmt.Println(msg.Text)
		MsgChan <- msg.Text
	}
}

// 測試案例
func Demo() {
	filename := `E:\xx.log`
	config := tail.Config{
		ReOpen:    true,
		Follow:    true,
		Location:  &tail.SeekInfo{Offset: 0, Whence: 2},
		MustExist: false,
		Poll:      true,
	}
	// 打開文件開始讀取數據
	tails, err := tail.TailFile(filename, config)
	if err != nil {
		fmt.Println("tails %s failed,err:%v\n", filename, err)
		return
	}
	var (
		msg *tail.Line
		ok  bool
	)
	fmt.Println("啟動")
	for {
		msg, ok = <-tails.Lines
		if !ok {
			fmt.Println("tails file close reopen,filename:$s\n", tails.Filename)
		}
		fmt.Println("msg:", msg.Text)
	}
}

到此這篇關于Golang監(jiān)聽日志文件并發(fā)送到kafka中的文章就介紹到這了,更多相關Golang 監(jiān)聽日志文件 內容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關文章希望大家以后多多支持腳本之家!

相關文章

  • Golang?websocket協(xié)議使用淺析

    Golang?websocket協(xié)議使用淺析

    這篇文章主要介紹了Golang?websocket協(xié)議的使用,WebSocket是一種新型的網絡通信協(xié)議,可以在Web應用程序中實現雙向通信,感興趣想要詳細了解可以參考下文
    2023-05-05
  • Go語言LeetCode500鍵盤行題解示例詳解

    Go語言LeetCode500鍵盤行題解示例詳解

    這篇文章主要為大家介紹了Go語言LeetCode500鍵盤行題解示例詳解,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進步,早日升職加薪
    2022-12-12
  • GO語言基礎入門第一個go程序解讀

    GO語言基礎入門第一個go程序解讀

    這篇文章主要為大家介紹了GO語言基礎入門的第一個go程序解讀,下面來帶大家進入Go語言世界helloworld的大門吧,有需要的朋友可以借鑒參考下,希望能夠有所幫助
    2021-11-11
  • VsCode搭建Go語言開發(fā)環(huán)境的配置教程

    VsCode搭建Go語言開發(fā)環(huán)境的配置教程

    這篇文章主要介紹了在VsCode中搭建Go開發(fā)環(huán)境的配置教程,本文給大家介紹的非常詳細,對大家的學習或工作具有一定的參考借鑒價值,需要的朋友可以參考下
    2020-05-05
  • Go語言使用select{}阻塞main函數介紹

    Go語言使用select{}阻塞main函數介紹

    這篇文章主要介紹了Go語言使用select{}阻塞main函數介紹,具有很好的參考價值,希望對大家有所幫助。一起跟隨小編過來看看吧
    2021-04-04
  • 深入理解Go中defer的機制

    深入理解Go中defer的機制

    本文主要介紹了Go中defer的機制,包括執(zhí)行順序、參數預計算、閉包和與返回值的交互,具有一定的參考價值,感興趣的可以了解一下
    2025-02-02
  • Go語言并發(fā)爬蟲的具體實現

    Go語言并發(fā)爬蟲的具體實現

    本文主要介紹了Go語言并發(fā)爬蟲的具體實現,文中通過示例代碼介紹的非常詳細,具有一定的參考價值,感興趣的小伙伴們可以參考一下
    2021-12-12
  • Go語言非main包編譯為靜態(tài)庫并使用的示例代碼

    Go語言非main包編譯為靜態(tài)庫并使用的示例代碼

    本文以Windows為例,介紹一下如何將Go的非main包編譯為靜態(tài)庫,用戶又將如何使用。通過實際項目創(chuàng)建常規(guī)工程,通過示例代碼給大家介紹的非常詳細,需要的朋友參考下吧
    2021-07-07
  • Golang之sync.Pool對象池對象重用機制總結

    Golang之sync.Pool對象池對象重用機制總結

    這篇文章主要對Golang的sync.Pool對象池對象重用機制做了一個總結,文中有相關的代碼示例和圖解,具有一定的參考價值,需要的朋友可以參考下
    2023-07-07
  • golang高性能的http請求 fasthttp詳解

    golang高性能的http請求 fasthttp詳解

    fasthttp 是 Go 的快速 HTTP 實現,當前在 1M 并發(fā)的生產環(huán)境使用非常成功,可以從單個服務器進行 100K qps 的持續(xù)連接,總而言之,fasthttp 比 net/http 快 10 倍,下面通過本文給大家介紹golang fasthttp http請求的相關知識,一起看看吧
    2021-09-09

最新評論

南安市| 孝感市| 安仁县| 台江县| 涟源市| 中宁县| 漳浦县| 平利县| 鄱阳县| 若尔盖县| 泸州市| 玉树县| 大新县| 大余县| 桂林市| 汉寿县| 林芝县| 高邮市| 吐鲁番市| 泰宁县| 深圳市| 静安区| 赤峰市| 新疆| 岱山县| 新巴尔虎左旗| 呼图壁县| 潜山县| 武威市| 三亚市| 高阳县| 普格县| 历史| 鄂州市| 区。| 葫芦岛市| 宽城| 射阳县| 江西省| 张家界市| 安平县|