【问题标题】:Can I use setWorkerCacheMb in Apache Beam 2.0+?我可以在 Apache Beam 2.0+ 中使用 setWorkerCacheMb 吗?
【发布时间】:2017-11-01 06:44:20
【问题描述】:

我的 Dataflow 作业(使用 Java SDK 2.1.0)非常慢,仅处理 50GB 就需要一天多的时间。我只是从 BigQuery (50GB) 中拉出一个完整的表,并从 GCS (100+MB) 中加入一个 csv 文件。

https://cloud.google.com/dataflow/model/group-by-key
我使用 sideInputs 执行连接(上面文档中的后一种方式),而我认为使用 CoGroupByKey 更有效,但我不确定这是我的工作超级慢的唯一原因。

我用谷歌搜索,默认情况下,sideinputs 的缓存设置为 100MB,我假设我的缓存稍微超过了这个限制,然后每个工作人员不断地重新读取 sideinputs。为了改进它,我想我可以使用setWorkerCacheMb 方法来增加缓存大小。

但是看起来DataflowPipelineOptions 没有这个方法并且DataflowWorkerHarnessOptions 被隐藏了。只需在 -Dexec.args 中传递 --workerCacheMb=200 就会导致

An exception occured while executing the Java class.
null: InvocationTargetException:
Class interface com.xxx.yyy.zzz$MyOptions missing a property
named 'workerCacheMb'. -> [Help 1]

如何使用此选项?谢谢。

我的管道:

MyOptions options = PipelineOptionsFactory.fromArgs(args).withValidation().as(MyOptions.class);

Pipeline p = Pipeline.create(options);

PCollection<TableRow> rows = p.apply("Read from BigQuery",
        BigQueryIO.read().from("project:MYDATA.events"));

// Read account file
PCollection<String> accounts = p.apply("Read from account file",
        TextIO.read().from("gs://my-bucket/accounts.csv")
                .withCompressionType(CompressionType.GZIP));
PCollection<TableRow> accountRows = accounts.apply("Convert to TableRow",
        ParDo.of(new DoFn<String, TableRow>() {
            private static final long serialVersionUID = 1L;

            @ProcessElement
            public void processElement(ProcessContext c) throws Exception {
                String line = c.element();
                CSVParser csvParser = new CSVParser();
                String[] fields = csvParser.parseLine(line);

                TableRow row = new TableRow();
                row = row.set("account_id", fields[0]).set("account_uid", fields[1]);
                c.output(row);
            }
        }));
PCollection<KV<String, TableRow>> kvAccounts = accountRows.apply("Populate account_uid:accounts KV",
        ParDo.of(new DoFn<TableRow, KV<String, TableRow>>() {
            private static final long serialVersionUID = 1L;

            @ProcessElement
            public void processElement(ProcessContext c) throws Exception {
                TableRow row = c.element();
                String uid = (String) row.get("account_uid");
                c.output(KV.of(uid, row));
            }
        }));
final PCollectionView<Map<String, TableRow>> uidAccountView = kvAccounts.apply(View.<String, TableRow>asMap());

// Add account_id from account_uid to event data
PCollection<TableRow> rowsWithAccountID = rows.apply("Join account_id",
        ParDo.of(new DoFn<TableRow, TableRow>() {
            private static final long serialVersionUID = 1L;

            @ProcessElement
            public void processElement(ProcessContext c) throws Exception {
                TableRow row = c.element();

                if (row.containsKey("account_uid") && row.get("account_uid") != null) {
                    String uid = (String) row.get("account_uid");
                    TableRow accRow = (TableRow) c.sideInput(uidAccountView).get(uid);
                    if (accRow == null) {
                        LOG.warn("accRow null, {}", row.toPrettyString());
                    } else {
                        row = row.set("account_id", accRow.get("account_id"));
                    }
                }
                c.output(row);
            }
        }).withSideInputs(uidAccountView));

// Insert into BigQuery
WriteResult result = rowsWithAccountID.apply(BigQueryIO.writeTableRows()
        .to(new TableRefPartition(StaticValueProvider.of("MYDATA"), StaticValueProvider.of("dev"),
                StaticValueProvider.of("deadletter_bucket")))
        .withFormatFunction(new SerializableFunction<TableRow, TableRow>() {
            private static final long serialVersionUID = 1L;

            @Override
            public TableRow apply(TableRow row) {
                return row;
            }
        }).withCreateDisposition(CreateDisposition.CREATE_NEVER)
        .withWriteDisposition(WriteDisposition.WRITE_APPEND));

p.run();

从历史上看,我的系统有两个用户标识符,新的 (account_id) 和旧的 (account_uid)。现在我需要将新的 account_id 添加到存储在 BigQuery 中的事件数据中,因为旧数据只有旧的 account_uid。 Accounts 表(account_uid 和 account_id 之间有关系)已经转换为 csv 并存储在 GCS 中。

最后一个函数TableRefPartition只是根据每个事件的时间戳将数据存储到BQ对应的分区中。作业仍在运行(2017-10-30_22_45_59-18169851018279768913),瓶颈看起来加入 account_id 部分。 这部分吞吐量(xxx 个元素/秒)根据图表上升和下降。根据图表,sideInputs 的估计大小为 106MB。

如果切换到 CoGroupByKey 可以显着提高性能,我会这样做。我只是懒惰,并认为使用 sideInputs 更容易处理没有帐户信息的事件数据。

【问题讨论】:

  • 请包含您的管道代码的代表性 sn-p,例如执行连接的 DoFn。我们或许可以帮助您找到通过其他方式对其进行优化的方法。
  • 谢谢@jkff,我添加了我的代码。我相信它只是遵循文档。请指出我是否做错了什么。

标签: google-cloud-dataflow apache-beam


【解决方案1】:

尝试以下之一:

1) 使用一些代码设置选项:

options.as(DataflowWorkerHarnessOptions.class).setWorkerCacheMb(500);

2) 让您的应用程序使用PipelineOptionsFactory 注册DataflowWorkerHarnessOptions

3) 拥有自己的选项类扩展 DataflowWorkerHarnessOptions

【讨论】:

  • 您好,第一个选项效果很好:)(我想其他人也会这样做)感谢您提供的信息!
  • 谢谢!这是否意味着 DWHO 的文档目前有些误导,应该进行更改以反映更改此选项实际上是安全的?
  • 怎么样? “这些选项在管道创建时没有影响。”表示这些选项仅适用于管道执行时。
【解决方案2】:

您可以采取一些措施来提高代码的性能:

  • 您的侧输入是Map&lt;String, TableRow&gt;,但您只使用了TableRow - accRow.get("account_id") 中的一个字段。将其设为Map&lt;String, String&gt; 怎么样,将其值设为account_id 本身?这可能比庞大的 TableRow 对象更有效。
  • 您可以将侧输入的值提取到 DoFn 中延迟初始化的成员变量中,以避免重复调用 .sideInput()

也就是说,这种表现是出乎意料的,我们正在调查是否还有其他事情发生。

【讨论】:

  • 我还看到这项工作在 GC 上花费了 很多 时间在工人身上。我的第一个建议肯定会对此有所帮助,但您也可以尝试指定更大的实例类型。
  • 谢谢,我改成 Map (account_id 实际上是一个 INTEGER) 类型,一个新工作超级快..(2017-10-31_15_33_41-2317555316151476602)。嗯,我了解到 TableRow 不是一个方便的数据结构来保存时间数据的好选择。您知道为什么在第一份工作中收集 TableRow 数据吗?我认为相同的 TableRow 实例会保留在每个工作人员的缓存中,直到工作完成。 “GC 运行”意味着数据被一次又一次地销毁和创建。是因为缓存大小(100Mb)的限制吗?
  • TableRow 是 BigQuery 特定的抽象,仅适用于两件事:从 BigQueryIO.read() 读取和写入 BigQueryIO.writeTableRows() :) 使用其他更多内容是个好主意管道中其他任何地方的高效类型。
  • 我怀疑这里的缓存确实保存了一堆 TableRow,但是每次缓存驱逐和随后的缓存未命中都会导致从 GCS 重新加载键/值存储的分片,并且访问磁盘非常昂贵,以至于您的作业性能受到这些罕见(因为您几乎适合缓存)但极其昂贵的事件的瓶颈。
  • 将侧输入的值提取到 DoFn 中延迟初始化的成员变量中是什么意思?我只在 processElement 函数中有上下文。你能举个例子吗?
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2017-01-29
  • 2019-01-25
相关资源
最近更新 更多