This is an automated email from the ASF dual-hosted git repository.
hansva pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/hop.git
The following commit(s) were added to refs/heads/main by this push:
new 64f9d190c7 fix db commit, fixes #8288 (#8480)
64f9d190c7 is described below
commit 64f9d190c77dca34f7657cfbe42a91508f4afe3e
Author: Hans Van Akelyen <[email protected]>
AuthorDate: Mon Sep 21 11:15:04 2026 +0200
fix db commit, fixes #8288 (#8480)
* fix db commit, fixes #8288
* extra hardening
---
.../monetdbbulkloader/MonetDbBulkLoader.java | 17 ++
.../monetdbbulkloader/MonetDbBulkLoaderTest.java | 39 +++
.../transforms/pgbulkloader/PGBulkLoader.java | 43 ++-
.../transforms/pgbulkloader/PGBulkLoaderTest.java | 100 +++++++
.../SynchronizeAfterMerge.java | 100 ++++++-
.../SynchronizeAfterMergeDisposeTest.java | 307 +++++++++++++++++++++
6 files changed, 596 insertions(+), 10 deletions(-)
diff --git
a/plugins/databases/monetdb/src/main/java/org/apache/hop/pipeline/transforms/monetdbbulkloader/MonetDbBulkLoader.java
b/plugins/databases/monetdb/src/main/java/org/apache/hop/pipeline/transforms/monetdbbulkloader/MonetDbBulkLoader.java
index 9ddb2d8caf..3c2d39e9ee 100644
---
a/plugins/databases/monetdb/src/main/java/org/apache/hop/pipeline/transforms/monetdbbulkloader/MonetDbBulkLoader.java
+++
b/plugins/databases/monetdb/src/main/java/org/apache/hop/pipeline/transforms/monetdbbulkloader/MonetDbBulkLoader.java
@@ -141,6 +141,23 @@ public class MonetDbBulkLoader extends
BaseTransform<MonetDbBulkLoaderMeta, Mone
return true;
}
+ /**
+ * The end-of-input branch of {@link #processRow()} flushes the buffer and
closes the MonetDB
+ * socket, but an error in the middle of the stream (the catch below) and a
stop that breaks the
+ * run loop mid-row never reach it - the socket, and the server-side load
session it holds, then
+ * stay open. Close it here as a backstop; {@link MapiSocket#close()} is
null-guarded and
+ * idempotent, so a normal, already-closed load is left untouched. See <a
+ * href="https://github.com/apache/hop/issues/8288">issue 8288</a>.
+ */
+ @Override
+ public void dispose() {
+ if (data.mserver != null) {
+ data.mserver.close();
+ data.mserver = null;
+ }
+ super.dispose();
+ }
+
@Override
public boolean processRow() throws HopException {
try {
diff --git
a/plugins/databases/monetdb/src/test/java/org/apache/hop/pipeline/transforms/monetdbbulkloader/MonetDbBulkLoaderTest.java
b/plugins/databases/monetdb/src/test/java/org/apache/hop/pipeline/transforms/monetdbbulkloader/MonetDbBulkLoaderTest.java
index 1d8c019797..c4f03f457e 100644
---
a/plugins/databases/monetdb/src/test/java/org/apache/hop/pipeline/transforms/monetdbbulkloader/MonetDbBulkLoaderTest.java
+++
b/plugins/databases/monetdb/src/test/java/org/apache/hop/pipeline/transforms/monetdbbulkloader/MonetDbBulkLoaderTest.java
@@ -20,6 +20,8 @@ package org.apache.hop.pipeline.transforms.monetdbbulkloader;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.verify;
import org.apache.hop.core.HopEnvironment;
import org.apache.hop.core.plugins.PluginRegistry;
@@ -30,6 +32,7 @@ import
org.apache.hop.pipeline.engines.local.LocalPipelineEngine;
import org.apache.hop.pipeline.transform.TransformMeta;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
+import org.monetdb.mcl.net.MapiSocket;
/** Test for MonetDbBulkLoader (excluding dialog). */
class MonetDbBulkLoaderTest {
@@ -84,6 +87,42 @@ class MonetDbBulkLoaderTest {
assertNull(MonetDbBulkLoader.hexFieldForMonetDbCopy(meta, null));
}
+ private MonetDbBulkLoader newLoader(MonetDbBulkLoaderData data) {
+ PipelineMeta pipelineMeta = new PipelineMeta();
+ TransformMeta transformMeta = new TransformMeta("test", new
MonetDbBulkLoaderMeta());
+ pipelineMeta.addTransform(transformMeta);
+ MonetDbBulkLoaderMeta meta = new MonetDbBulkLoaderMeta();
+ Pipeline pipeline = new LocalPipelineEngine(pipelineMeta);
+ return new MonetDbBulkLoader(transformMeta, meta, data, 0, pipelineMeta,
pipeline);
+ }
+
+ /**
+ * An error or a stop in the middle of the stream skips the end-of-input
close. dispose() has to
+ * close the MonetDB socket so the load session does not stay open on the
server. Issue 8288.
+ */
+ @Test
+ void disposeClosesTheOpenSocket() {
+ MonetDbBulkLoaderData data = new MonetDbBulkLoaderData();
+ MapiSocket mserver = mock(MapiSocket.class);
+ data.mserver = mserver;
+
+ newLoader(data).dispose();
+
+ verify(mserver).close();
+ assertNull(data.mserver);
+ }
+
+ /** A transform that never opened a socket must dispose cleanly. */
+ @Test
+ void disposeSurvivesWithoutASocket() {
+ MonetDbBulkLoaderData data = new MonetDbBulkLoaderData();
+ data.mserver = null;
+
+ newLoader(data).dispose();
+
+ assertNull(data.mserver);
+ }
+
@Test
void testEscapeOsPathSpacesOnWindows() {
PipelineMeta pipelineMeta = new PipelineMeta();
diff --git
a/plugins/databases/postgresql/src/main/java/org/apache/hop/pipeline/transforms/pgbulkloader/PGBulkLoader.java
b/plugins/databases/postgresql/src/main/java/org/apache/hop/pipeline/transforms/pgbulkloader/PGBulkLoader.java
index 16a4291888..68b06877f2 100644
---
a/plugins/databases/postgresql/src/main/java/org/apache/hop/pipeline/transforms/pgbulkloader/PGBulkLoader.java
+++
b/plugins/databases/postgresql/src/main/java/org/apache/hop/pipeline/transforms/pgbulkloader/PGBulkLoader.java
@@ -26,6 +26,7 @@ package org.apache.hop.pipeline.transforms.pgbulkloader;
//
import com.google.common.annotations.VisibleForTesting;
+import java.io.IOException;
import java.math.BigDecimal;
import java.nio.charset.Charset;
import java.sql.Connection;
@@ -229,7 +230,11 @@ public class PGBulkLoader extends
BaseTransform<PGBulkLoaderMeta, PGBulkLoaderDa
pgCopyOut.flush();
pgCopyOut.endCopy();
pgCopyOut.close();
- data.db.getConnection().close();
+ pgCopyOut = null;
+ }
+ if (data != null && data.db != null) {
+ data.db.disconnect();
+ data.db = null;
}
return false;
@@ -481,6 +486,42 @@ public class PGBulkLoader extends
BaseTransform<PGBulkLoaderMeta, PGBulkLoaderDa
}
}
+ /**
+ * The end-of-input branch of {@link #processRow()} finishes the COPY and
closes the connection. A
+ * stop or an error in the middle of the stream never reaches it, leaving
the COPY and its
+ * connection open on the server - and with them the locks the load holds.
Abort the copy here
+ * ({@code endCopy()} would commit the partial rows, {@code cancelCopy()}
discards them) and
+ * release the connection. See <a
href="https://github.com/apache/hop/issues/8288">issue 8288</a>.
+ */
+ @Override
+ public void dispose() {
+ try {
+ if (pgCopyOut != null && pgCopyOut.isActive()) {
+ pgCopyOut.cancelCopy();
+ }
+ } catch (SQLException e) {
+ logError("Error cancelling the COPY command while stopping the
transform", e);
+ } finally {
+ try {
+ // Only close a copy that is no longer active. A still-active copy
here means cancelCopy()
+ // above threw (a broken connection, the likely case), and pgjdbc's
close() runs endCopy()
+ // on an active copy - which would commit the very rows we are trying
to discard. The
+ // disconnect() below tears the connection down regardless.
+ if (pgCopyOut != null && !pgCopyOut.isActive()) {
+ pgCopyOut.close();
+ }
+ } catch (IOException e) {
+ logError("Error closing the COPY output stream", e);
+ }
+ pgCopyOut = null;
+ if (data.db != null) {
+ data.db.disconnect();
+ data.db = null;
+ }
+ }
+ super.dispose();
+ }
+
protected void verifyDatabaseConnection() throws HopException {
// Confirming Database Connection is defined.
if (meta.getConnection() == null) {
diff --git
a/plugins/databases/postgresql/src/test/java/org/apache/hop/pipeline/transforms/pgbulkloader/PGBulkLoaderTest.java
b/plugins/databases/postgresql/src/test/java/org/apache/hop/pipeline/transforms/pgbulkloader/PGBulkLoaderTest.java
index 716fab3dbe..598a1b3d97 100644
---
a/plugins/databases/postgresql/src/test/java/org/apache/hop/pipeline/transforms/pgbulkloader/PGBulkLoaderTest.java
+++
b/plugins/databases/postgresql/src/test/java/org/apache/hop/pipeline/transforms/pgbulkloader/PGBulkLoaderTest.java
@@ -27,12 +27,16 @@ import static org.junit.jupiter.api.Assertions.fail;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.Mockito.doNothing;
import static org.mockito.Mockito.doReturn;
+import static org.mockito.Mockito.doThrow;
import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.never;
import static org.mockito.Mockito.spy;
import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.when;
+import java.lang.reflect.Field;
import java.nio.charset.StandardCharsets;
+import java.sql.SQLException;
import java.util.ArrayList;
import org.apache.hop.core.HopClientEnvironment;
import org.apache.hop.core.database.Database;
@@ -54,6 +58,7 @@ import org.junit.jupiter.api.BeforeAll;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.RegisterExtension;
+import org.postgresql.copy.PGCopyOutputStream;
class PGBulkLoaderTest {
@RegisterExtension
@@ -246,6 +251,101 @@ class PGBulkLoaderTest {
}
}
+ private PGBulkLoader disposableLoader(PGBulkLoaderData data,
PGCopyOutputStream copyOut)
+ throws Exception {
+ PGBulkLoader loader =
+ spy(
+ new PGBulkLoader(
+ transformMockHelper.transformMeta,
+ transformMockHelper.iTransformMeta,
+ data,
+ 0,
+ transformMockHelper.pipelineMeta,
+ transformMockHelper.pipeline));
+ Field field = PGBulkLoader.class.getDeclaredField("pgCopyOut");
+ field.setAccessible(true);
+ field.set(loader, copyOut);
+ return loader;
+ }
+
+ /**
+ * A stop or an error leaves the COPY open. dispose() must abort it - not
endCopy(), which would
+ * commit the partial rows - and release the connection. Issue 8288.
+ */
+ @Test
+ void disposeCancelsAnActiveCopyAndDisconnects() throws Exception {
+ PGBulkLoaderData data = new PGBulkLoaderData();
+ Database db = mock(Database.class);
+ data.db = db;
+ PGCopyOutputStream copyOut = mock(PGCopyOutputStream.class);
+ // Active on entry (so the copy is cancelled), inactive afterwards (a
successful cancel), so the
+ // stream is then closed without an endCopy() commit.
+ when(copyOut.isActive()).thenReturn(true, false);
+
+ PGBulkLoader loader = disposableLoader(data, copyOut);
+ loader.dispose();
+
+ verify(copyOut).cancelCopy();
+ verify(copyOut, never()).endCopy();
+ verify(copyOut).close();
+ verify(db).disconnect();
+ assertNull(data.db);
+ }
+
+ /**
+ * The normal end-of-input path already finished and closed the COPY.
dispose() then only has to
+ * release the connection, and must not touch the already-closed copy.
+ */
+ @Test
+ void disposeDoesNotCancelAnAlreadyFinishedCopy() throws Exception {
+ PGBulkLoaderData data = new PGBulkLoaderData();
+ Database db = mock(Database.class);
+ data.db = db;
+ PGCopyOutputStream copyOut = mock(PGCopyOutputStream.class);
+ when(copyOut.isActive()).thenReturn(false);
+
+ PGBulkLoader loader = disposableLoader(data, copyOut);
+ loader.dispose();
+
+ verify(copyOut, never()).cancelCopy();
+ verify(copyOut).close();
+ verify(db).disconnect();
+ }
+
+ /**
+ * If cancelCopy() throws (a broken connection mid-load, the likely case),
dispose() must not fall
+ * through to close() - pgjdbc's close() runs endCopy() on a still-active
copy, committing the
+ * very rows we are discarding. The connection is torn down instead. Issue
8288 review follow-up.
+ */
+ @Test
+ void disposeDoesNotCommitWhenCancelFailsOnAnActiveCopy() throws Exception {
+ PGBulkLoaderData data = new PGBulkLoaderData();
+ Database db = mock(Database.class);
+ data.db = db;
+ PGCopyOutputStream copyOut = mock(PGCopyOutputStream.class);
+ when(copyOut.isActive()).thenReturn(true);
+ doThrow(new SQLException("connection reset")).when(copyOut).cancelCopy();
+
+ PGBulkLoader loader = disposableLoader(data, copyOut);
+ loader.dispose();
+
+ verify(copyOut).cancelCopy();
+ verify(copyOut, never()).endCopy();
+ verify(copyOut, never()).close();
+ verify(db).disconnect();
+ assertNull(data.db);
+ }
+
+ /** A transform stopped before the first row opened neither the copy nor the
connection. */
+ @Test
+ void disposeSurvivesWithoutACopyOrConnection() throws Exception {
+ PGBulkLoaderData data = new PGBulkLoaderData();
+ data.db = null;
+ PGBulkLoader loader = disposableLoader(data, null);
+
+ loader.dispose();
+ }
+
private static PGBulkLoaderMeta getPgBulkLoaderMock(String DbNameOverride)
throws HopXmlException {
PGBulkLoaderMeta pgBulkLoaderMetaMock = mock(PGBulkLoaderMeta.class);
diff --git
a/plugins/transforms/synchronizeaftermerge/src/main/java/org/apache/hop/pipeline/transforms/synchronizeaftermerge/SynchronizeAfterMerge.java
b/plugins/transforms/synchronizeaftermerge/src/main/java/org/apache/hop/pipeline/transforms/synchronizeaftermerge/SynchronizeAfterMerge.java
index 70e1a17b68..5b2506aae3 100644
---
a/plugins/transforms/synchronizeaftermerge/src/main/java/org/apache/hop/pipeline/transforms/synchronizeaftermerge/SynchronizeAfterMerge.java
+++
b/plugins/transforms/synchronizeaftermerge/src/main/java/org/apache/hop/pipeline/transforms/synchronizeaftermerge/SynchronizeAfterMerge.java
@@ -1083,10 +1083,74 @@ public class SynchronizeAfterMerge
return databaseMeta;
}
+ /**
+ * A single-threaded (streaming) pipeline never sends the end-of-input
signal that {@link
+ * #processRow()} flushes on; it calls this after every batch of rows
instead. Commit what is
+ * pending now, so it does not sit uncommitted - and, on databases like
Oracle, locked - until the
+ * stream ends. See <a
href="https://github.com/apache/hop/issues/8288">issue 8288</a>.
+ */
+ @Override
+ public void batchComplete() throws HopException {
+ if (data.db == null || data.db.getConnection() == null) {
+ return;
+ }
+ try {
+ emptyBatchBuffer(false);
+ } catch (HopDatabaseBatchException be) {
+ // The statements stay open for the next batch, so drop what failed
before going on. The rest
+ // is the recovery a failure in the middle of the stream gets.
+ for (PreparedStatement statement : data.preparedStatements.values()) {
+ data.db.clearBatch(statement);
+ }
+ if (getTransformMeta().isDoingErrorHandling()) {
+ data.db.commit(true);
+ processBatchException(be.toString(), be.getUpdateCounts(),
be.getExceptionsList());
+ } else {
+ data.db.rollback();
+ throw new HopException(
+ BaseMessages.getString(PKG,
"SynchronizeAfterMerge.Error.UpdatingBatch"), be);
+ }
+ } catch (SQLException e) {
+ throw new HopDatabaseException("Unexpected error committing the database
connection.", e);
+ }
+ }
+
+ /**
+ * The end-of-input path in {@link #processRow()} has normally flushed and
disconnected by now. It
+ * is skipped when the transform is stopped or fails in the middle of a row,
and a single-threaded
+ * (streaming) pipeline never sends end-of-input at all. In both cases the
connection stayed open,
+ * and with it the uncommitted transaction and every row lock it holds. See
<a
+ * href="https://github.com/apache/hop/issues/8288">issue 8288</a>.
+ *
+ * <p>A graceful stop leaves {@code getErrors() == 0}, so the pending batch
is committed rather
+ * than rolled back - deliberately, and in line with Table Output, Update
and Delete: this
+ * transform already commits every {@code commitSize} rows, so committing
the final partial batch
+ * on a stop keeps the same all-or-a-multiple-of-commitSize contract. Only a
real error rolls
+ * back. Note that {@link #emptyBatchBuffer(boolean)} still calls {@code
putRow} for the committed
+ * rows; on a stop {@code putRow} is a no-op (nothing reads downstream
anyway), while the rows are
+ * safely in the table - the same behaviour Table Output has.
+ */
+ @Override
+ public void dispose() {
+ if (data.db != null) {
+ if (data.db.getConnection() != null) {
+ if (getErrors() > 0) {
+ // The transform failed: nothing that is still pending may reach the
table.
+ rollback();
+ data.db.disconnect();
+ } else {
+ finishTransform();
+ }
+ }
+ data.db = null;
+ }
+ super.dispose();
+ }
+
private void finishTransform() {
if (data.db != null && data.db.getConnection() != null) {
try {
- finishTransformEmptyBatchBuffer();
+ emptyBatchBuffer(true);
} catch (HopDatabaseBatchException be) {
finishTransformErrorHandling(be);
} catch (Exception dbe) {
@@ -1098,11 +1162,7 @@ public class SynchronizeAfterMerge
setOutputDone();
if (getErrors() > 0) {
- try {
- data.db.rollback();
- } catch (HopDatabaseException e) {
- logError("Unexpected error rolling back the database connection.",
e);
- }
+ rollback();
}
data.db.disconnect();
@@ -1110,7 +1170,21 @@ public class SynchronizeAfterMerge
}
}
- private void finishTransformEmptyBatchBuffer()
+ private void rollback() {
+ try {
+ data.db.rollback();
+ } catch (HopDatabaseException e) {
+ logError("Unexpected error rolling back the database connection.", e);
+ }
+ }
+
+ /**
+ * Execute and commit the pending batch of every prepared statement and pass
the buffered rows on.
+ *
+ * @param closeStatements true at the end of the transform; false between
the batches of a
+ * single-threaded pipeline, where the statements are reused.
+ */
+ private void emptyBatchBuffer(boolean closeStatements)
throws SQLException, HopDatabaseException, HopTransformException,
HopValueException {
if (!data.db.getConnection().isClosed()) {
for (String schemaTable : data.preparedStatements.keySet()) {
@@ -1121,9 +1195,17 @@ public class SynchronizeAfterMerge
batchCounter = 0;
}
- PreparedStatement insertStatement =
data.preparedStatements.get(schemaTable);
+ // Between batches there is nothing to commit for a statement that
took no rows this batch.
+ // Skip it; at final completion we still fall through so
emptyAndCommit closes the
+ // statement.
+ if (!closeStatements && batchCounter == 0) {
+ continue;
+ }
+
+ PreparedStatement statement = data.preparedStatements.get(schemaTable);
- data.db.emptyAndCommit(insertStatement, data.batchMode, batchCounter);
+ data.db.emptyAndCommit(statement, data.batchMode, batchCounter,
closeStatements);
+ data.commitCounterMap.put(schemaTable, 0);
}
for (int i = 0; i < data.batchBuffer.size(); i++) {
Object[] row = data.batchBuffer.get(i);
diff --git
a/plugins/transforms/synchronizeaftermerge/src/test/java/org/apache/hop/pipeline/transforms/synchronizeaftermerge/SynchronizeAfterMergeDisposeTest.java
b/plugins/transforms/synchronizeaftermerge/src/test/java/org/apache/hop/pipeline/transforms/synchronizeaftermerge/SynchronizeAfterMergeDisposeTest.java
new file mode 100644
index 0000000000..e5cc731a25
--- /dev/null
+++
b/plugins/transforms/synchronizeaftermerge/src/test/java/org/apache/hop/pipeline/transforms/synchronizeaftermerge/SynchronizeAfterMergeDisposeTest.java
@@ -0,0 +1,307 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.hop.pipeline.transforms.synchronizeaftermerge;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.ArgumentMatchers.anyBoolean;
+import static org.mockito.ArgumentMatchers.anyInt;
+import static org.mockito.ArgumentMatchers.anyLong;
+import static org.mockito.ArgumentMatchers.anyString;
+import static org.mockito.ArgumentMatchers.eq;
+import static org.mockito.ArgumentMatchers.nullable;
+import static org.mockito.Mockito.doAnswer;
+import static org.mockito.Mockito.doNothing;
+import static org.mockito.Mockito.doReturn;
+import static org.mockito.Mockito.doThrow;
+import static org.mockito.Mockito.inOrder;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.never;
+import static org.mockito.Mockito.spy;
+import static org.mockito.Mockito.verify;
+
+import java.sql.BatchUpdateException;
+import java.sql.Connection;
+import java.sql.PreparedStatement;
+import java.util.ArrayList;
+import java.util.List;
+import org.apache.hop.core.database.Database;
+import org.apache.hop.core.exception.HopDatabaseBatchException;
+import org.apache.hop.core.exception.HopException;
+import org.apache.hop.core.row.IRowMeta;
+import org.apache.hop.core.row.RowMeta;
+import org.apache.hop.core.row.value.ValueMetaString;
+import org.apache.hop.pipeline.PipelineMeta;
+import org.apache.hop.pipeline.engines.local.LocalPipelineEngine;
+import org.apache.hop.pipeline.transform.TransformMeta;
+import org.apache.hop.pipeline.transform.TransformPartitioningMeta;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+import org.mockito.InOrder;
+
+/**
+ * Releasing the database connection when the end of the input never comes.
+ *
+ * <p>The transform used to flush and disconnect only from the end-of-input
branch of {@code
+ * processRow()}. A stop or a failure in the middle of a row skips that
branch, and a
+ * single-threaded (streaming) pipeline never reaches it at all: the
connection then stayed open for
+ * as long as the JVM lived, and with it the uncommitted transaction and the
row locks it held on
+ * the database. See <a href="https://github.com/apache/hop/issues/8288">issue
8288</a>.
+ *
+ * <p>The invariants pinned here: {@code dispose()} always ends with {@code
disconnect()} when a
+ * connection is open, commits pending work unless the transform failed, and
{@code batchComplete()}
+ * commits between batches without closing the statements the next batch
reuses.
+ */
+class SynchronizeAfterMergeDisposeTest {
+
+ private static final String INSERT_KEY = "\"T\"" +
SynchronizeAfterMerge.CONST_INSERT;
+ private static final String UPDATE_KEY = "\"T\"" +
SynchronizeAfterMerge.CONST_UPDATE;
+
+ private SynchronizeAfterMerge transform;
+ private SynchronizeAfterMergeData data;
+ private TransformMeta transformMeta;
+ private Database db;
+ private PreparedStatement insertStatement;
+ private PreparedStatement updateStatement;
+ private List<Object[]> emitted;
+ private List<String> rejected;
+
+ @BeforeEach
+ void setUp() throws Exception {
+ SynchronizeAfterMergeMeta meta = mock(SynchronizeAfterMergeMeta.class);
+ transformMeta = mock(TransformMeta.class);
+ doReturn("transform").when(transformMeta).getName();
+ doReturn(mock(TransformPartitioningMeta.class))
+ .when(transformMeta)
+ .getTargetTransformPartitioningMeta();
+ doReturn(meta).when(transformMeta).getTransform();
+
+ PipelineMeta pipelineMeta = mock(PipelineMeta.class);
+ doReturn(transformMeta).when(pipelineMeta).findTransform(anyString());
+
+ db = mock(Database.class);
+ Connection connection = mock(Connection.class);
+ doReturn(connection).when(db).getConnection();
+
+ insertStatement = mock(PreparedStatement.class);
+ updateStatement = mock(PreparedStatement.class);
+
+ data = new SynchronizeAfterMergeData();
+ data.db = db;
+ data.batchMode = true;
+ data.insertValue = "insert";
+ data.indexOfOperationOrderField = 1;
+ data.inputRowMeta = new RowMeta();
+ data.inputRowMeta.addValueMeta(new ValueMetaString("name"));
+ data.inputRowMeta.addValueMeta(new ValueMetaString("operation"));
+ data.outputRowMeta = data.inputRowMeta;
+ data.preparedStatements.put(INSERT_KEY, insertStatement);
+ data.preparedStatements.put(UPDATE_KEY, updateStatement);
+ data.commitCounterMap.put(INSERT_KEY, 3);
+ data.commitCounterMap.put(UPDATE_KEY, 2);
+ data.batchBuffer = new ArrayList<>();
+
+ transform =
+ spy(
+ new SynchronizeAfterMerge(
+ transformMeta, meta, data, 1, pipelineMeta, spy(new
LocalPipelineEngine())));
+ doReturn(transformMeta).when(transform).getTransformMeta();
+ doReturn(false).when(transform).isRowLevel();
+ doNothing().when(transform).logDetailed(anyString());
+ doNothing().when(transform).logError(anyString());
+ doNothing().when(transform).logError(anyString(), any(Throwable.class));
+
+ emitted = new ArrayList<>();
+ rejected = new ArrayList<>();
+ doAnswer(
+ inv -> {
+ emitted.add(inv.getArgument(1));
+ return null;
+ })
+ .when(transform)
+ .putRow(any(IRowMeta.class), any());
+ doAnswer(
+ inv -> {
+ rejected.add(inv.getArgument(3));
+ return null;
+ })
+ .when(transform)
+ .putError(
+ any(IRowMeta.class),
+ any(),
+ anyLong(),
+ anyString(),
+ nullable(String.class),
+ anyString());
+ }
+
+ private void bufferRows(int count) {
+ for (int i = 0; i < count; i++) {
+ data.batchBuffer.add(new Object[] {"row" + i, "insert"});
+ }
+ }
+
+ /** The stop-in-the-middle-of-a-row and the streaming case: nothing flushed
us before. */
+ @Test
+ void disposeFlushesCommitsAndDisconnectsWhenEndOfInputNeverCame() throws
Exception {
+ bufferRows(5);
+
+ transform.dispose();
+
+ // Both statements are flushed and closed. Their relative order is a
Hashtable artifact, so it
+ // is
+ // not asserted; what matters is that a flush precedes the disconnect.
+ verify(db).emptyAndCommit(updateStatement, true, 2, true);
+ InOrder order = inOrder(db);
+ order.verify(db).emptyAndCommit(insertStatement, true, 3, true);
+ order.verify(db).disconnect();
+ verify(db, never()).rollback();
+ assertEquals(5, emitted.size(), "the buffered rows leave on the output
before we disconnect");
+ assertTrue(data.batchBuffer.isEmpty());
+ assertNull(data.db, "the database handle is released for the garbage
collector");
+ }
+
+ /** After a failure nothing that is still pending may reach the table. */
+ @Test
+ void disposeRollsBackAndDisconnectsAfterAFailure() throws Exception {
+ bufferRows(5);
+ transform.setErrors(1);
+
+ transform.dispose();
+
+ InOrder order = inOrder(db);
+ order.verify(db).rollback();
+ order.verify(db).disconnect();
+ verify(db, never()).emptyAndCommit(any(), anyBoolean(), anyInt(),
anyBoolean());
+ assertTrue(emitted.isEmpty(), "no row is reported as written after a
rollback");
+ }
+
+ /** The normal end-of-input path already disconnected: dispose() must not do
it twice. */
+ @Test
+ void disposeIsANoOpOnceTheConnectionIsGone() throws Exception {
+ doReturn(null).when(db).getConnection();
+
+ transform.dispose();
+
+ verify(db, never()).disconnect();
+ verify(db, never()).rollback();
+ assertNull(data.db);
+ }
+
+ /** A transform whose init() never got as far as a connection. */
+ @Test
+ void disposeSurvivesAMissingDatabase() throws Exception {
+ data.db = null;
+
+ transform.dispose();
+
+ verify(db, never()).disconnect();
+ }
+
+ /** Between batches the statements are reused: commit, but leave them open.
*/
+ @Test
+ void batchCompleteCommitsWithoutClosingTheStatements() throws Exception {
+ bufferRows(5);
+
+ transform.batchComplete();
+
+ verify(db).emptyAndCommit(insertStatement, true, 3, false);
+ verify(db).emptyAndCommit(updateStatement, true, 2, false);
+ verify(db, never()).disconnect();
+ assertEquals(0, data.commitCounterMap.get(INSERT_KEY), "the batch counter
starts over");
+ assertEquals(0, data.commitCounterMap.get(UPDATE_KEY), "the batch counter
starts over");
+ assertEquals(5, emitted.size());
+ assertTrue(data.batchBuffer.isEmpty());
+ }
+
+ /**
+ * Between batches a statement that took no rows must not trigger a no-op
commit; only the
+ * statements with pending work are flushed. Issue 8288 review follow-up.
+ */
+ @Test
+ void batchCompleteSkipsStatementsWithAnEmptyBatch() throws Exception {
+ data.commitCounterMap.put(UPDATE_KEY, 0);
+ bufferRows(5);
+
+ transform.batchComplete();
+
+ verify(db).emptyAndCommit(insertStatement, true, 3, false);
+ verify(db, never()).emptyAndCommit(eq(updateStatement), anyBoolean(),
anyInt(), anyBoolean());
+ }
+
+ @Test
+ void batchCompleteIsANoOpWithoutAConnection() throws Exception {
+ doReturn(null).when(db).getConnection();
+
+ transform.batchComplete();
+
+ verify(db, never()).emptyAndCommit(any(), anyBoolean(), anyInt(),
anyBoolean());
+ }
+
+ /**
+ * A batch that fails between batches gets the same recovery as one that
fails mid-stream: the
+ * failed batches are dropped so the next batch does not re-run them, what
went through is
+ * committed and every buffered row leaves on one stream or the other.
+ */
+ @Test
+ void batchCompleteWithErrorHandlingRoutesTheFailedBatchAndKeepsGoing()
throws Exception {
+ doReturn(true).when(transformMeta).isDoingErrorHandling();
+ bufferRows(5);
+ BatchUpdateException cause =
+ new BatchUpdateException("boom", new int[] {1, 1,
java.sql.Statement.EXECUTE_FAILED});
+ HopDatabaseBatchException failure =
+ Database.createHopDatabaseBatchException("Error updating batch",
cause);
+ doThrow(failure)
+ .when(db)
+ .emptyAndCommit(eq(insertStatement), anyBoolean(), anyInt(),
eq(false));
+
+ transform.batchComplete();
+
+ verify(db).clearBatch(insertStatement);
+ verify(db).clearBatch(updateStatement);
+ verify(db).commit(true);
+ verify(db, never()).rollback();
+ verify(db, never()).disconnect();
+ assertEquals(5, emitted.size() + rejected.size(), "every buffered row
leaves exactly once");
+ assertTrue(data.batchBuffer.isEmpty());
+ }
+
+ /** Without error handling a failed batch is a failed transform. */
+ @Test
+ void batchCompleteWithoutErrorHandlingRollsBackAndFails() throws Exception {
+ doReturn(false).when(transformMeta).isDoingErrorHandling();
+ bufferRows(5);
+ HopDatabaseBatchException failure =
+ Database.createHopDatabaseBatchException(
+ "Error updating batch", new BatchUpdateException("boom", new
int[0]));
+ doThrow(failure)
+ .when(db)
+ .emptyAndCommit(eq(insertStatement), anyBoolean(), anyInt(),
eq(false));
+
+ assertThrows(HopException.class, () -> transform.batchComplete());
+
+ verify(db).clearBatch(insertStatement);
+ verify(db).clearBatch(updateStatement);
+ verify(db).rollback();
+ verify(db, never()).commit(anyBoolean());
+ verify(db, never()).disconnect();
+ }
+}