【问题标题】:Using Spanner within Apache Beam Dataflow在 Apache Beam 数据流中使用 Spanner
【发布时间】:2018-12-27 07:54:54
【问题描述】:

我正在尝试在 Apache Beam ParDo(DoFn) 中添加 Spanner 连接。我需要查找一些行作为 ParDo 的一部分。数据流创建了许多工作人员(通常最多 4 个),我使用 startBundle 和 finishBundle 方法在工作人员的生命周期内打开和关闭扳手连接。然后在 processElement 方法中,我对传递 DatabaseClient 并使用 singleUseReadOnlyTransaction 的每个项目执行查找。

我应该添加它作为 GCP 下的数据流运行

一些代码来说明这一点。

private static CustomDoFn<String, TransactionImport> processRow = new CustomDoFn<String, TransactionImport>(){
    private static final long serialVersionUID = 1L;

    private Spanner spanner = null;
    private DatabaseClient dbClient = null;

    @StartBundle
    public void startBundle(StartBundleContext c){
      TransactionFileOptions options = c.getPipelineOptions().as(TransactionFileOptions.class);

      com.google.cloud.spanner.SpannerOptions spannerOptions = com.google.cloud.spanner.SpannerOptions.newBuilder().build();
      spanner = spannerOptions.getService();
      String spannerProjectID = options.getSpannerProjectId();
      String spannerInstanceID = options.getSpannerInstanceId();
      String spannerDatabaseID = options.getSpannerDatabaseId();

      DatabaseId db = DatabaseId.of(spannerProjectID, spannerInstanceID, spannerDatabaseID);
      dbClient = spanner.getDatabaseClient(db);
    }

    @FinishBundle
    public void finishBundle(FinishBundleContext c){
        spanner.close();  
    }

    @ProcessElement
    public void processElement(DoFn<String, TransactionImport>.ProcessContext c) throws Exception {
    TransactionImport import = new TransactionImport();

    Statement statement = Statement.newBuilder("SELECT * FROM Table1 WHERE Name= @Name")
            .bind("Name").to( text)
            .build();

    ResultSet resultSet = dbClient.singleUseReadOnlyTransaction().executeQuery(statement);

    // set some value  on import dependant on retrieved value

    c.output(import);

}

这总是导致数据流未完成,当我检查日志时,我看到:

Processing stuck in step Process Rows for at least 05m00s without outputting or completing in state process
at sun.misc.Unsafe.park(Native Method)
at java.util.concurrent.locks.LockSupport.park(LockSupport.java:175)
at java.util.concurrent.SynchronousQueue$TransferStack.awaitFulfill(SynchronousQueue.java:458)
at java.util.concurrent.SynchronousQueue$TransferStack.transfer(SynchronousQueue.java:362)
at java.util.concurrent.SynchronousQueue.take(SynchronousQueue.java:924)
at com.google.common.util.concurrent.Uninterruptibles.takeUninterruptibly(Uninterruptibles.java:233)
at com.google.cloud.spanner.SessionPool$Waiter.take(SessionPool.java:411)
at com.google.cloud.spanner.SessionPool$Waiter.access$3300(SessionPool.java:399)
at com.google.cloud.spanner.SessionPool.getReadSession(SessionPool.java:754)
at com.google.cloud.spanner.DatabaseClientImpl.singleUseReadOnlyTransaction(DatabaseClientImpl.java:52)
at com.mycompany.pt.SpannerDataAccess.getBinDetails(SpannerDataAccess.java:197)
at com.mycompany.pt.transactionFiles.TransactionFileDataflow$1.processLine(TransactionFileDataflow.java:411)
at com.mycompany.pt.transactionFiles.TransactionFileDataflow$1.processElement(TransactionFileDataflow.java:336)
at com.mycompany.pt.transactionFiles.TransactionFileDataflow$1$DoFnInvoker.invokeProcessElement(Unknown Source)

`

有人在 ParDo 中使用过这样的 Spanner 吗?

【问题讨论】:

  • 我在处理大批量并尝试将数据插入 BigQuery 时遇到了同样的问题。我会尝试使用 stream_insert 而不是 batch_insert 并让您知道它是否有效。

标签: google-cloud-platform google-cloud-dataflow apache-beam google-cloud-spanner


【解决方案1】:

我不是扳手专家,但也许我可以提供帮助:

  1. 您应该使用@Setup/@Teardown 连接和断开扳手。 @{Start,Finish}Bundle 在工作人员的生命周期内被多次调用。更多详情请看这里:https://beam.apache.org/documentation/execution-model/#bundling-and-persistence

  2. 您的 processElement 方法是否曾经使用 c.output(...)?如果没有,Beam 会认为你的管道卡住了

【讨论】:

  • 感谢 Igor,使用 Setup 和 Teardown 方法的问题在于它们不使用上下文参数来传递我的 Spanner 参数。我的 processElement 确实包含一个 c.output,我只是不想在 sn-p 中放入太多代码。再次感谢
  • @RichardB 请在 sn-p 中包含c.output 以避免混淆。
  • c.output 添加到原始帖子以进行澄清
  • @RichardB 这就是 SpannerIO 需要数据库/实例 ID 作为转换参数的原因之一。作为替代方案,您可以通过实例 ID 和数据库 ID 对 PCollection 中的元素进行分组,然后自己进行捆绑。并在单个@ProcessElement 中建立连接并运行事务
猜你喜欢
  • 1970-01-01
  • 2018-07-20
  • 2019-06-04
  • 2021-05-20
  • 1970-01-01
  • 2019-09-16
  • 1970-01-01
  • 1970-01-01
  • 2019-04-23
相关资源
最近更新 更多