Как применить параллельное программирование golang к реальной разработке

задняя часть

Предисловие: Я писал инструментальный скрипт для анализа большого количества лог-файлов онлайн, это должно было быть скучной работой, но принцип достижения максимума вдохновил меня постоянно думать о том, как его оптимизировать. В этой статье будет немного объяснен процесс оптимизации из первой версии в процессе разработки, и, наконец, использование golang для реализации пула рабочих потоков, похожего на java, который полон вознаграждений.

1. Безмозглая открытая стадия горутины

1. Предыстория задачи Функция этого инструмента кратко представлена ​​следующим образом: во-первых, количество онлайн-журналов очень велико, затем необходимо прочитать содержимое файла журнала, а затем элементы журнала анализируются один за другим, и нужный журнал элементы записываются в файл, а затем извлекаются.Данные в этом файле вычисляются. Данные имитации формата записи журнала следующие:

{
 "id":xx,"time":"2017-11-19","key1":"value1","key2":"value2".......  
}

2. Безмозглый горунтин Сначала дайте код (поместите только основной код, опустите операции с файлами и обработку ошибок исключений), а затем расскажите об идее этого этапа.

var wirteChan = make(chan []byte) //用于写入文件
var waitgroup sync.WaitGroup //用于控制同步
func main(){
    //省略写入文件的打开操作,毕竟我们主要讲并发这块
    
    //初始化写入的channel
	InitWriter(outLogFileWriter)
	//接下来去遍历每个日志文件,每读出一个日志项就开一个gorountine去处理
	for _,file := range logDir {
	    if file.IsDir() {
			continue
		} else {
			HandlerFile(arg.dir + "/" + file.Name()) //处理每个文件
		}
	}
}
/**
 * 初始化Writer的channel
 */
func InitWriter(outLogFileWriter *bufio.Writer) {
	go func() {
		for data := range wirteChan {
			nn, err := outLogFileWriter.Write(data)
		}
	}()
}
//处理每个文件,然后开G去处理每个日志项
func HandlerFile(fileName string) {
	file, err := os.Open(fileName)
	defer file.Close()
	br := bufio.NewReader(file)
	for {
		data, err := br.ReadBytes('\n')
		if err == io.EOF {
			break
		} else {
			go Handler(data) //每次开一个G去处理,处理完写入writeChannel
		}
	}
}

Анализ: писать в блоге не нравится код увеличения длины, поэтому важно, чтобы чуть выше кода над кодом комментировали, когда дело доходит до того, чтобы не повторяться. Что ж, давайте подумаем, по приведенному выше коду есть вопросы? Теперь мы предполагаем, что у нас есть только очень большой файл, файл для каждой записи записи не имеет мозгов, чтобы открыть G для выполнения. Затем запустите его, мы обнаружим, что он очень медленный ~. В чем проблема? Во-первых, мы не можем контролировать количество G, а затем файлы журнала очень велики, поэтому бежать вниз, количество G очень велико, несколько G для записи данных в канал, тогда произойдет серьезная блокировка. По разным причинам вел к этому методу неприменимо

Во-вторых, присоединиться к очереди задач с буферизацией

1. Очередь задач Выше мы сказали, что мы не можем контролировать количество задач, поэтому я добавил сюда очередь задач, чтобы ставить задачи в очередь и одновременно контролировать количество задач. Над кодом:

/**
 * Job结构体,包含要处理的数据和处理函数(这个可根据需要修改)
 */
type Job struct {
	Data []byte
	Proc func([]byte)
}
//Job队列,存储要做的Job,将每个任务打包成Job发送到这里
var JobQueue chan Job = make(chan Job, arg.maxqueue)
//启动处理函数处理
func Handler(Data []byte) {
    for range job := <-Queue {
        job.Proc(Data)
    }
}

Анализ: в настоящее время задание модели задачи абстрагировано.Поскольку вызов функции на самом деле представляет собой адрес функции плюс параметры функции, мы также можем поместить функцию обработки в задание. Затем позвольте функции-обработчику обработать это. Думая об этом, я немного любуюсь собой, а потом с большим интересом запускаю ее. Ну, это не кажется намного быстрее (на самом деле это зависит от вашей функции обработки, которая является Proc в Job). какие? Успокойся и проанализируй это, я действительно думаю, что я очень милый. Я просто обернул задачу, а затем использовал буферизованную очередь задач.Поскольку созданное задание намного больше, чем вычислительная мощность одного M, буферизация просто немного задержала проблему.

В-третьих, модель Работа/Рабочий

На самом деле, когда я пишу здесь, у меня уже есть немного B-числа о том, как оптимизировать. Я вспомнил концепцию пула потоков в java, я могу построить пул потоков, и тогда пул содержит несколько воркеров (количество можно указать), каждый воркер идет в очередь, чтобы получить задачи для обработки, и продолжает получать задачи после обработка. В то же время, чтобы повысить универсальность, типы параметров изменены на interface{}. Хорошо, давайте посмотрим на код дальше, код здесь очень критичен, поэтому я все это выложил

type Job struct {
	Data interface{}
	Proc func(interface{})
}
//Job队列,存储要做的Job
var JobQueue chan Job = make(chan Job, arg.maxqueue)
//Woker,用来从Job队列中取出Job执行
type Worker struct {
	WokerPool  chan chan Job //表示属于哪个Worker池,同时接收JobChannel注册
	JobChannel chan Job      //任务管道,通过这个管道获取任务执行
	Quit       chan bool     //用来停止Worker
}
//新建一个Worker,需要传入Worker池参数
func NewWorker(wokerPool chan chan Job) Worker {
	return Worker{
		WokerPool:  wokerPool,
		JobChannel: make(chan Job),
		Quit:       make(chan bool),
	}
}
//Worker的启动:包含:(1) 把该worker的JobChannel注册到WorkerPool中去  (2) 监听JobChannel上有没有新的任务到来 (3) 监听是否受到关闭的请求
func (worker Worker) Start() {
	go func() {
		for {
			worker.WokerPool <- worker.JobChannel //每次做完任务后就重新注册上去通知本worker又处于可用状态了
			select {
			case job := <-worker.JobChannel:
				job.Proc(job.Data)
			case quit := <-worker.Quit: //接收到关闭信息,直接退出即可
				if quit {
					return
				}
			}
		}
	}()
}
//Worker的关闭:只要发送一个关闭信号即可
func (worker Worker) Stop() {
	go func() {
		worker.Quit <- true
	}()
}
//管理Worker的调度器,包含最大worker数量和workerpool
type Dispatcher struct {
	MaxWorker  int
	WorkerPool chan chan Job
}

//启动一个调度器
func (dispatcher *Dispatcher) Run() {
	//启动maxworker个worker
	for i := 0; i < dispatcher.MaxWorker; i++ {
		worker := NewWorker(dispatcher.WorkerPool)
		worker.Start()
	}
	//接下来启动调度服务
	go dispatcher.dispatch()
}

func (dispatcher *Dispatcher) dispatch() {
	for {
		select {
		case job := <-JobQueue:
			go func(job Job) {
				jobChannel := <-dispatcher.WorkerPool //获取一个可用的worker
				jobChannel <- job                     //将该job发送给该worker
			}(job)
		}
	}
}

//新建一个调度器
func NewDispatcher(maxWorker int) *Dispatcher {
	workerPool := make(chan chan Job, maxWorker)
	return &Dispatcher{
		WorkerPool: workerPool,
		MaxWorker:  maxWorker,
	}
}

Анализ: Каждое предложение в коде очень четко аннотировано и не будет повторяться. Мы можем запустить модель следующим образом:dispatcher := NewDispatcher(MaxWorker) dispatcher.Run(). Важно подчеркнуть, что это лечение необходимо написать функцию в соответствии с их собственным бизнесом, а затем упаковываться в данные, а затем распределяться на работу на работу на линии. Затем я запускаю свой скрипт, десятки функций обработки файлов G в три раунда (то есть, что мне нужны три лечения, каждый раунд лечения основан на результатах последнего раунда), занял от трех минут до четырех минут, а размещение процессоров Скорость не высока. Для высокой трудоемки вы можете использовать инструмент для анализа его в конце PPROF медленного в то, где

В-четвертых, резюме

Потому что я только что изучил принцип параллелизма в golang раньше, а потом у меня как раз была эта задача, поэтому я начал исследовать и оптимизировать с нуля.После написания всего инструмента у меня есть более глубокое понимание параллелизма в golang.Блокировки, файловые операции и т. д. также знакомы. Я приобрел много вещей, поэтому я призываю изучать новое, не только понимать принцип, но и делать больше самостоятельно, чтобы быть твердым. На самом деле в этой модели все еще есть некоторые недостатки, которые в дальнейшем будут оптимизироваться. В этот период я ​​также упомянул несколько очень хороших блогов, и я хотел бы выразить здесь свою благодарность.