【问题标题】:Inserting rows on BigQuery: InsertAllRequest Vs BigQueryIO.writeTableRows()在 BigQuery 上插入行:InsertAllRequest 与 BigQueryIO.writeTableRows()
【发布时间】:2018-12-21 07:01:16
【问题描述】:

当我使用 writeTableRows 在 BigQuery 上插入行时,与 InsertAllRequest 相比,性能确实很差。显然,有些东西没有正确设置。

用例 1: 我编写了一个 Java 程序来使用 Twitter4j 处理“示例”Twitter 流。当一条推文出现时,我使用以下命令将其写入 BigQuery:

insertAllRequestBuilder.addRow(rowContent);

当我从我的 Mac 运行这个程序时,它每分钟将大约 1000 行直接插入 BigQuery 表中。我认为通过在集群上运行 Dataflow 作业可以做得更好。

用例 2: 当一条推文进来时,我将它写到 Google 的 PubSub 的主题中。我从我的 Mac 上运行它,它每分钟发送大约 1000 条消息。

我编写了一个 Dataflow 作业,该作业读取此主题并使用 BigQueryIO.writeTableRows() 写入 BigQuery。我有一个 8 机器 Dataproc 集群。我使用 DataflowRunner 在该集群的主节点上开始了这项工作。它慢得令人难以置信!就像每 5 分钟左右 100 行一样。这是相关代码的sn-p:

statuses.apply("ToBQRow", ParDo.of(new DoFn<Status, TableRow>() {
    @ProcessElement
    public void processElement(ProcessContext c) throws Exception {
        TableRow row = new TableRow();
        Status status = c.element();
        row.set("Id", status.getId());
        row.set("Text", status.getText());
        row.set("RetweetCount", status.getRetweetCount());
        row.set("FavoriteCount", status.getFavoriteCount());
        row.set("Language", status.getLang());
        row.set("ReceivedAt", null);
        row.set("UserId", status.getUser().getId());
        row.set("CountryCode", status.getPlace().getCountryCode());
        row.set("Country", status.getPlace().getCountry());
        c.output(row);
    }
})) 
    .apply("WriteTableRows", BigQueryIO.writeTableRows().to(tweetsTable)//
            .withSchema(schema)
            .withMethod(BigQueryIO.Write.Method.FILE_LOADS)
            .withTriggeringFrequency(org.joda.time.Duration.standardMinutes(2))
            .withNumFileShards(1000)
            .withWriteDisposition(BigQueryIO.Write.WriteDisposition.WRITE_APPEND)
            .withCreateDisposition(BigQueryIO.Write.CreateDisposition.CREATE_IF_NEEDED));

我做错了什么?我应该使用“SparkRunner”吗?如何确认它在我的集群的所有节点上运行?

【问题讨论】:

  • 您能否明确说明 Dataproc 如何参与您的用例。如果您使用的是 Dataflow 运行器,这将启动一些 GCE 虚拟机(工作者)来运行该作业。您是否尝试为 Cloud Pipeline 更改 parameters?您可以设置更多的 numWorkers 并更改 workerMachineType。
  • 我的错! DataflowRunner 将在托管模式下运行。我的帐户不允许我使用超过 4 名工人,因此速度提升并不显着。从文档中不清楚我需要在哪个服务中请求增加配额。如果您知道,请告诉我。我也会继续寻找。感谢您的帮助。
  • 你应该增加Compute Engine API CPUs的配额

标签: google-cloud-platform google-bigquery google-cloud-pubsub google-cloud-dataproc dataflow


【解决方案1】:

使用 BigQuery,您可以:

  • 流式传输数据。低延迟(高达每秒 10 万行)是有成本的。
  • 批量输入数据。延迟更高,吞吐量惊人,完全免费。

这就是你正在经历的不同。如果您只想摄取 1000 行,批处理会明显变慢。对于 100 亿行,通过批处理将更快,而且无需任何成本。

Dataflow/Bem 的 BigQueryIO.writeTableRows 可以流式传输或批处理数据。

BigQueryIO.Write.Method.FILE_LOADS 粘贴的代码选择批处理。

【讨论】:

  • 当我将其更改为 BigQueryIO.Write.Method.STREAMING_INSERTS 时,它的性能更好,但总体速度仍然很慢。有趣的是,“ToBQRow”步骤非常慢,这没有任何意义,因为它所做的只是创建一个新的 TableRow 并编写它。有什么办法可以加快速度吗?
  • ToBQRow 开始计数。输入集合 -> 添加的元素 -> 13,829。输出集合 -> 添加元素 -> 249. 哇...这一步太慢了
猜你喜欢
  • 1970-01-01
  • 2018-06-12
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2019-02-19
  • 1970-01-01
相关资源
最近更新 更多