【问题标题】:Go routine for processing files in a work flow. What am I missing?在工作流程中处理文件的例行程序。我错过了什么?
【发布时间】:2021-11-19 14:25:12
【问题描述】:

我正在构建一些东西来监视文件上传的目录。现在我正在使用 for {} 循环不断读取目录以进行测试,并计划在未来使用 cron 或其他东西来启动我的应用程序。

目标是监控上传目录,确保文件已完成复制,然后将文件移动到另一个目录进行处理。这些文件本身的大小从 15GB 到大约 50GB 不等,我们每天会收到数百个。

这是我第一次涉足围棋程序。我不确定我是否完全误解了 go 例程、通道和等待组或其他东西,但我认为当我遍历文件列表时,每个文件都会由 goroutine 函数独立处理。但是,当我运行下面的代码时,它会抓取一个文件,但只确认它在目录中找到的第一个文件。我注意到,一旦第一个文件完成,其他文件就会被确认为已完成。

package main

import (
        "flag"
        "fmt"
        "io/ioutil"
        "log"
        "os"
        "sync"
        "time"

        "gopkg.in/yaml.v2"
)

type Config struct {
        LogFileName  string `yaml:"logfilename"`
        LogFilePath  string `yaml:"logfilepath"`
        UploadRoot   string `yaml:"upload_root"`
        TPUploadTool string `yaml:"tp_upload_tool"`
}

var wg sync.WaitGroup

const WORKERS = 5

func getConfig(fileName string) (*Config, error) {
        conf := &Config{}
        yamlFile, err := os.Open(fileName)
        if err != nil {
                fmt.Printf("Error reading YAML file: %s\n", err)
                os.Exit(1)
        }

        defer yamlFile.Close()
        yaml_decoder := yaml.NewDecoder(yamlFile)

        if err := yaml_decoder.Decode(conf); err != nil {
                return nil, err
        }
        return conf, err
}

func getFileData(fileToUpload string, fileStatus chan string) {
        var newSize int64
        var currentSize int64

        currentSize = 0
        newSize = 0
        fmt.Printf("Uploading: %s\n", fileToUpload)
        fileDone := false
        for !fileDone {
                fileToUploadStat, _ := os.Stat(fileToUpload)
                currentSize = fileToUploadStat.Size()
                //fmt.Printf("%s current size is: %d\n", fileToUpload, currentSize)
                //fmt.Println("New size ", newSize)
                if currentSize != 0 {
                        if currentSize > newSize {
                                newSize = currentSize
                        } else if newSize == currentSize {
                                fileStatus <- "Done"
                                fileDone = true
                                wg.Done()
                        }
                }
                time.Sleep(1 * time.Second)
        }

}

func sendToCDS() {
        fmt.Println("Sending To CDS")
}

func main() {

        fileStatus := make(chan string)
        configFileName := flag.String("config", "", "YAML configuration file.\n")
        flag.Parse()

        if *configFileName == "" {
                flag.PrintDefaults()
                os.Exit(1)
        }

        UploaderConfig, err := getConfig(*configFileName)
        if err != nil {
                log.Fatal("Error reading configuration file.")
        }

        for {
                fmt.Print("Checking for new files..")
                uploadFiles, err := ioutil.ReadDir(UploaderConfig.UploadRoot)
                if err != nil {
                        log.Fatal(err)
                }
                if len(uploadFiles) == 0 {
                        fmt.Println("..no files to transfer.\n")
                }
                for _, uploadFile := range uploadFiles {
                        wg.Add(1)
                        fmt.Println("...adding", uploadFile.Name())
                        if err != nil {
                                log.Fatalln("Unable to read file information.")
                        }
                        ff := UploaderConfig.UploadRoot + "/" + uploadFile.Name()

                        go getFileData(ff, fileStatus)
                        status := <-fileStatus

                        if status == "Done" {
                                fmt.Printf("%s is done.\n", uploadFile.Name())
                                os.Remove(ff)
                        }
                }
                wg.Wait()
        }
}

我曾考虑将通道用于线程安全排队机制,该机制加载目录中的文件,然后文件被工作人员拾取。我在 Python 中做过类似的事情。

【问题讨论】:

  • 我觉得不用channel,wait group就够了。
  • 是的,我看到人们使用频道作为工作/工作队列来做类似于我正在尝试做的事情的不同示例。我不确定等待小组是否会完成同样的事情。

标签: go


【解决方案1】:

由于以下几行,代码按顺序处理每个文件:

go getFileData(ff, fileStatus)
status := <-fileStatus

第一行创建了一个 goroutine,但第二行一直等到该 goroutine 完成它的工作。

如果你想并行处理文件,那么你可以使用工作池模式。

jobs:=make(chan string)
done:=make(chan struct{})
for i:=0;i<nWorkers;i++ {
   go workerFunc(jobs,done)
}

jobs 频道将用于将新发现的文件发送给工作人员。当您发现一个新文件时,您可以简单地:

jobs <- fileName

并且工作人员应该处理文件,并且应该返回从通道读取。所以它应该是这样的:

func worker(ch chan string,done chan struct{}) {
   defer func() {
      done<-struct{}{} // Notify that this goroutine is completed
   }()
   for inputFile:=range ch {
       // process inputFile
   }
}

当一切都完成后,您可以通过等待所有 goroutine 完成来终止程序:

close(jobs)
for i:=0;i<nWorkers;i++ {
   <-done
}

【讨论】:

  • 这就是我最初的目标。等待组会是更好的选择吗?
  • 等待组将替换完成的通道。这两种方案在功能上是等价的。
猜你喜欢
  • 2013-03-15
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多