【发布时间】: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