【问题标题】:Writing data from bigquery to csv is slow将数据从 bigquery 写入 csv 很慢
【发布时间】:2022-01-20 10:35:13
【问题描述】:

我编写的代码表现得很奇怪而且很慢,我不明白为什么。 我要做的是将数据从 bigquery(使用查询作为输入)下载到 CSV 文件,然后使用此 CSV 创建一个 url 链接,以便人们可以将其作为报告下载。 我正在尝试优化编写 CSV 的过程,因为它需要一些时间并且有一些奇怪的行为。

代码迭代 bigquery 结果并将每个结果传递到通道以供将来使用 golang encoding/csv 包进行解析/编写。 这是一些调试的相关部分

func (s *Service) generateReportWorker(ctx context.Context, query, reportName string) error {
    it, err := s.bigqueryClient.Read(ctx, query)
    if err != nil {
        return err
    }
    filename := generateReportFilename(reportName)
    gcsObj := s.gcsClient.Bucket(s.config.GcsBucket).Object(filename)
    wc := gcsObj.NewWriter(ctx)
    wc.ContentType = "text/csv"
    wc.ContentDisposition = "attachment"

    csvWriter := csv.NewWriter(wc)

    var doneCount uint64

    go backgroundTimer(ctx, it.TotalRows, &doneCount)

    rowJobs := make(chan []bigquery.Value, it.TotalRows)
    workers := 10
    wg := sync.WaitGroup{}
    wg.Add(workers)

    // start wrokers pool
    for i := 0; i < workers; i++ {
        go func(c context.Context, num int) {
            defer wg.Done()
            for row := range rowJobs {
                records := make([]string, len(row))
                for j, r := range records {
                    records[j] = fmt.Sprintf("%v", r)
                }
                s.mu.Lock()
                start := time.Now()
                if err := csvWriter.Write(records); err != {
                    log.Errorf("Error writing row: %v", err)
                }
                if time.Since(start) > time.Second {
                    fmt.Printf("worker %d took %v\n", num, time.Since(start))
                }
                s.mu.Unlock()
                atomic.AddUint64(&doneCount, 1)
            }
        }(ctx, i)
    }

    // read results from bigquery and add to the pool
    for {
        var row []bigquery.Value
        if err := it.Next(&row); err != nil {
            if err == iterator.Done || err == context.DeadlineExceeded {
                break
            }
            log.Errorf("Error loading next row from BQ: %v", err)
        }
        rowJobs <- row
    }

    fmt.Println("***done loop!***")

    close(rowJobs)

    wg.Wait()

    csvWriter.Flush()
    wc.Close()

    url := fmt.Sprintf("%s/%s/%s", s.config.BaseURL s.config.GcsBucket, filename)

    /// ....

}

func backgroundTimer(ctx context.Context, total uint64, done *uint64) {
    ticker := time.NewTicker(10 * time.Second)
    go func() {
        for {
            select {
            case <-ctx.Done():
                ticker.Stop()
                return
            case _ = <-ticker.C:
                fmt.Printf("progress (%d,%d)\n", atomic.LoadUint64(done), total)
            }
        }
    }()
}

bigquery 读取函数

func (c *Client) Read(ctx context.Context, query string) (*bigquery.RowIterator, error)  {
    job, err := c.bigqueryClient.Query(query).Run(ctx)
    if err != nil {
        return nil, err
    }
    it, err := job.Read(ctx)
    if err != nil {
        return nil, err
    }
    return it, nil
}

我使用大约 400,000 行的查询运行此代码。查询本身大约需要 10 秒,但整个过程大约需要 2 分钟 输出:

progress (112346,392565)
progress (123631,392565)
***done loop!***
progress (123631,392565)
progress (123631,392565)
progress (123631,392565)
progress (123631,392565)
progress (123631,392565)
progress (123631,392565)
progress (123631,392565)
worker 3 took 1m16.728143875s
progress (247525,392565)
progress (247525,392565)
progress (247525,392565)
progress (247525,392565)
progress (247525,392565)
progress (247525,392565)
progress (247525,392565)
worker 3 took 1m13.525662666s
progress (370737,392565)
progress (370737,392565)
progress (370737,392565)
progress (370737,392565)
progress (370737,392565)
progress (370737,392565)
progress (370737,392565)
progress (370737,392565)
worker 4 took 1m17.576536375s
progress (392565,392565)

您可以看到写入前 112346 行的速度很快,然后由于某种原因,worker 3 花了 1.16 分钟 (!!!) 来写入一行,这导致其他 worker 等待互斥锁被释放,而这又发生了 2 次,导致整个过程需要 2 多分钟才能完成。

我不确定发生了什么,我该如何进一步调试,为什么我在执行中会出现这个停顿?

【问题讨论】:

  • 用本地文件替换gcsObj.NewWriter(ctx)需要多少时间?
  • 我认为这就是问题所在!写入本地文件而不是云存储要快得多(整个过程需要 30 秒).. 在云存储中,由于某种原因,在写入 ~123631 行之后写入下一个字节非常慢(源代码行是 b.Flush() in bufio.去,第 710 行,在函数 WriteString())
  • 这样您就可以将所有记录写入本地文件,然后一次调用将文件传输到 GCP。然后你就不需要多个工人了。
  • 现在 io.Copy 从本地到 gcs 很慢
  • 如果用gcloud命令行工具复制速度一样吗?

标签: csv go google-bigquery


【解决方案1】:

按照@serge-v 的建议,您可以将所有记录写入本地文件,然后将文件作为一个整体传输到 GCS。为了使该过程在更短的时间内发生,您可以将文件分成多个块并可以使用以下命令:gsutil -m cp -j where

gsutil 用于从命令行访问云存储

-m用于执行并行多线程/多处理复制

cp用于复制文件

-j 将 gzip 传输编码应用于任何文件上传。这也节省了网络带宽,同时将未压缩的数据保留在 Cloud Storage 中。

要在你的 go 程序中应用这个命令,你可以参考这个Github link.

您可以尝试在您的 Go 程序中实现 profiling。分析将帮助您分析复杂性。也可以通过profiling找到程序中的耗时。

由于您正在从 BigQuery 读取数百万行,因此您可以尝试使用 BigQuery Storage API。与批量数据导出相比,它提供对 BigQuery 托管存储的更快访问。使用 BigQuery Storage API 而不是您在 Go 程序中使用的迭代器可以加快处理速度。

如需更多参考,您还可以查看 BigQuery 提供的 Query Optimization techniques

【讨论】:

  • 我查看了 Bigquery Storage API,但在我的情况下导致性能低下的不是结果的迭代,而是将结果写入文件。也许使用它会使过程更快,但不值得开销(我不希望超过 2M 行)
  • 嗨@Avishay28,我已经用 serge-v 建议的方式更新了答案,因为它可以帮助您解决问题,以及获得更快性能的其他替代方法,即使用 BigQuery Storage API .
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2021-12-10
  • 2021-05-17
  • 1970-01-01
  • 2017-12-03
  • 1970-01-01
  • 1970-01-01
  • 2014-10-04
相关资源
最近更新 更多