【问题标题】:Write to Postgres with apache beam (GCP)使用 apache beam (GCP) 写入 Postgres
【发布时间】:2021-10-10 06:08:36
【问题描述】:

我们在我们的谷歌云平台中使用 apache Beam,并实现了一个数据流流作业,该作业写入我们的 postgres 数据库。然而,我们注意到,一旦我们开始使用两个相邻的 JdbcIO.write() 语句,我们的流式作业开始抛出如下错误:

Operation ongoing in step JdbcIO.WriteVoid/ParDo(Write) for at least 35m00s without outputting or completing in state process
at jdk.internal.misc.Unsafe.park (Native Method)
at java.util.concurrent.locks.LockSupport.park (LockSupport.java:194)
at java.util.concurrent.locks.AbstractQueuedSynchronizer$ConditionObject.await (AbstractQueuedSynchronizer.java:2081)
at org.apache.commons.pool2.impl.LinkedBlockingDeque.takeFirst (LinkedBlockingDeque.java:581)
at org.apache.commons.pool2.impl.GenericObjectPool.borrowObject (GenericObjectPool.java:439)
at org.apache.commons.pool2.impl.GenericObjectPool.borrowObject (GenericObjectPool.java:356)
at org.apache.commons.dbcp2.PoolingDataSource.getConnection (PoolingDataSource.java:134)
at org.apache.beam.sdk.io.jdbc.JdbcIO$WriteVoid$WriteFn.executeBatch (JdbcIO.java:1438)
at org.apache.beam.sdk.io.jdbc.JdbcIO$WriteVoid$WriteFn.processElement (JdbcIO.java:1387)
at org.apache.beam.sdk.io.jdbc.JdbcIO$WriteVoid$WriteFn$DoFnInvoker.invokeProcessElement (Unknown Source)

这仅在部署后大约 30 分钟发生。它能够处理 10.000 个元素,直到 30 分钟后。平均吞吐量范围从 50 个元素/秒到 120 个元素/秒。 查询也不是那么繁重,只是一个简单的delete 和insert 语句。 我们认为其他元素的连接被卡住并且没有释放,但我们不知道如何修复它。

代码如下:

public void writeToPostgres(PCollection<TimestampedValue<KV<String, Duration>>> collection) {
  collection
    .apply(Filter.by(Postgres::filter1))
    .apply(JdbcIO.<TimestampedValue<KV<String, Duration>>>write()
      .withDataSourceProviderFn(JdbcIO.PoolableDataSourceProvider.of(getDataSourceConfiguration()))
      .withStatement("DELETE FROM table1 where field1 = ?::UUID and field2=?")
      .withPreparedStatementSetter((element, statement) -> {
        statement.setString(1, element.getValue().getKey());
        Instant timestamp = element.getTimestamp();
        statement.setTimestamp(2, new Timestamp(timestamp.getMillis()));
      })
    .withBatchSize(1)
    .withRetryStrategy(DEADLOCK_DETECTED_RETRY_STRATEGY));

  collection
    .apply(Filter.by(Postgres::filter2))
    .apply(
      JdbcIO.<TimestampedValue<KV<String, Duration>>>write()
        .withDataSourceProviderFn(JdbcIO.PoolableDataSourceProvider.of(getDataSourceConfiguration()))
    .withStatement("INSERT INTO table1 (field1, field2) \n" +
      "VALUES (?::UUID, ?) \n" +
      "ON CONFLICT ON CONSTRAINT someconstraint\n" +
      "DO UPDATE SET field2 = excluded.field2")
    .withPreparedStatementSetter((element, statement) -> {
      Instant eventTime = element.getTimestamp();
      Timestamp now = Timestamp.from(now());

      statement.setString(1, element.getValue().getKey());
      statement.setTimestamp(2, new Timestamp(eventTime.getMillis()));
    })
    .withBatchSize(1)
    .withRetryStrategy(DEADLOCK_DETECTED_RETRY_STRATEGY)
  );
}

...

private DataSourceConfiguration getDataSourceConfiguration() {
    return DataSourceConfiguration.create(ValueProvider.StaticValueProvider.of("org.postgresql.Driver"), jdbcUrlProvider)
                .withUsername(usernameProvider)
                .withPassword(passwordProvider);
}

我该如何解决这个问题?

【问题讨论】:

  • 我对 JdbcIO 知之甚少,无法将其作为完整答案,但通读 the code 看起来你是对的,在 ParDo 最终确定之前,连接是打开的,而不是关闭的,这发生在管道完成时。这两个转换是否都连接到同一个数据源?
  • @DanielOliveira 是的,这是正确的。 update 和insert 语句在同一个数据库和表上执行,用户名和密码相同。我将编辑问题以显示getDataSourceConfiguration 函数的外观。

标签: java postgresql google-cloud-platform google-cloud-dataflow apache-beam


【解决方案1】:

我们能够找到解决方法,但我认为这更像是一种解决方法,因为我们在 JdbcIO 的 DataSourceProvider 中没有找到任何东西。我们基本上复制了JdbcIO 的PoolableDataSourceProvider 并改用了HikariDataSource,因为无论如何它似乎可以提高性能。

首先,我们在 out pom 文件中添加 hikariCP 依赖

<dependency>
     <groupId>com.zaxxer</groupId>
     <artifactId>HikariCP</artifactId>
     <version>5.0.0</version>
</dependency>

HikariDataSourceProvider 如下所示:

public static class HikariDataSourceProvider implements SerializableFunction<Void, DataSource> {
    private static final ConcurrentHashMap<HikariDataSourceConfig, DataSource> instances = new ConcurrentHashMap<>();

    private final HikariDataSourceConfig config;

    private HikariDataSourceProvider(HikariDataSourceConfig config) {
        this.config = config;
    }

    public static SerializableFunction<Void, DataSource> of(HikariDataSourceConfig hikariDataSourceConfig) {
        return new HikariDataSourceProvider(hikariDataSourceConfig);
    }

    @Override
    public DataSource apply(Void input) {
        return instances.computeIfAbsent(
                config,
                ignored -> {
                    HikariDataSource hikariDataSource = new HikariDataSource();
                        hikariDataSource.setJdbcUrl(config.getJdbcUrlProvider().get());
                        hikariDataSource.setUsername(config.getUsernameProvider().get());
                        hikariDataSource.setPassword(config.getPasswordProvider().get());
                        hikariDataSource.setAutoCommit(false);
                        return hikariDataSource;
                });
    }
}
...

@Data
@Builder
public static class HikariDataSourceConfig implements Serializable {
    private final ValueProvider<String> jdbcUrlProvider;
    private final ValueProvider<String> usernameProvider;
    private final ValueProvider<String> passwordProvider;
}

@Data 和 @Builder 是 lombok 注释。

PTransform 看起来像这样:

JdbcIO.<TimestampedValue<KV<String, Duration>>>write()
     .withDataSourceProviderFn(HikariDataSourceProvider.of(getDataSourceConfig()))
     .withStatement("...

我们还删除了.withBatchSize(1) 行,因此它不会成为流程的瓶颈。我们尝试在没有 HikariDataSource 的情况下先删除此行,但仅此一项并不能解决此问题。

流式作业现在可以处理语句并且很稳定。错误不再发生。

【讨论】:

    猜你喜欢
    • 2021-03-25
    • 2021-05-20
    • 2020-05-01
    • 1970-01-01
    • 2021-04-29
    • 2020-06-23
    • 2020-12-21
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多