PostgresChannelMessageTableSubscriber: Renew connection only if invalid
Fixes: #9111 An evolution of the #9061: renew the connection only when we need to. **Auto-cherry-pick to `6.2.x` & `6.1.x`**
This commit is contained in:
committed by
Artem Bilan
parent
899598a342
commit
da29e2da6a
@@ -216,9 +216,10 @@ public final class PostgresChannelMessageTableSubscriber implements SmartLifecyc
|
||||
if (!isActive()) {
|
||||
return;
|
||||
}
|
||||
if (notifications == null || notifications.length == 0) {
|
||||
if ((notifications == null || notifications.length == 0) && !conn.isValid(1)) {
|
||||
//We did not receive any notifications within the timeout period.
|
||||
//We will close the connection and re-establish it.
|
||||
//If the connection is still valid, we will continue polling
|
||||
//Otherwise, we will close the connection and re-establish it.
|
||||
break;
|
||||
}
|
||||
for (PGNotification notification : notifications) {
|
||||
|
||||
@@ -23,8 +23,10 @@ import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
import java.util.concurrent.CountDownLatch;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.concurrent.atomic.AtomicBoolean;
|
||||
import java.util.concurrent.atomic.AtomicInteger;
|
||||
import java.util.concurrent.atomic.AtomicReference;
|
||||
import java.util.function.Consumer;
|
||||
|
||||
import javax.sql.DataSource;
|
||||
|
||||
@@ -268,7 +270,18 @@ public class PostgresChannelMessageTableSubscriberTests implements PostgresConta
|
||||
CountDownLatch latch = new CountDownLatch(2);
|
||||
List<Object> payloads = new ArrayList<>();
|
||||
CountDownLatch connectionLatch = new CountDownLatch(2);
|
||||
connectionSupplier.onGetConnection = connectionLatch::countDown;
|
||||
AtomicBoolean connectionCloseState = new AtomicBoolean();
|
||||
connectionSupplier.onGetConnection = conn -> {
|
||||
connectionLatch.countDown();
|
||||
if (connectionCloseState.compareAndSet(false, true)) {
|
||||
try {
|
||||
conn.close();
|
||||
}
|
||||
catch (Exception e) {
|
||||
//nop
|
||||
}
|
||||
}
|
||||
};
|
||||
postgresChannelMessageTableSubscriber.start();
|
||||
postgresSubscribableChannel.subscribe(message -> {
|
||||
payloads.add(message.getPayload());
|
||||
@@ -324,7 +337,7 @@ public class PostgresChannelMessageTableSubscriberTests implements PostgresConta
|
||||
|
||||
private static class ConnectionSupplier implements PgConnectionSupplier {
|
||||
|
||||
Runnable onGetConnection;
|
||||
Consumer<PgConnection> onGetConnection;
|
||||
|
||||
@Override
|
||||
public PgConnection get() throws SQLException {
|
||||
@@ -333,10 +346,11 @@ public class PostgresChannelMessageTableSubscriberTests implements PostgresConta
|
||||
POSTGRES_CONTAINER.getPassword())
|
||||
.unwrap(PgConnection.class);
|
||||
if (this.onGetConnection != null) {
|
||||
this.onGetConnection.run();
|
||||
this.onGetConnection.accept(conn);
|
||||
}
|
||||
return conn;
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user