MartijnVisser commented on code in PR #246:
URL: 
https://github.com/apache/flink-connector-jdbc/pull/246#discussion_r4135488649


##########
flink-connector-jdbc-core/src/test/java/org/apache/flink/connector/jdbc/core/table/sink/JdbcOutputFormatTest.java:
##########
@@ -588,6 +599,327 @@ void testInvalidConnectionInJdbcOutputFormat() throws 
IOException, SQLException
         }
     }
 
+    @Test
+    void testUpsertBranchWithNativeUpsertReducesByKey() throws Exception {
+        RecordingDialect dialect = new RecordingDialect(true);
+
+        assertChangelogIsReducedByKey(dialect);
+
+        // the dialect's own upsert is used; the insert-or-update fallback is 
never assembled
+        assertThat(dialect.upsertCalls).isEqualTo(1);
+        assertThat(dialect.deleteCalls).isEqualTo(1);
+        assertThat(dialect.rowExistsCalls).isZero();
+        assertThat(dialect.updateCalls).isZero();
+    }
+
+    @Test
+    void testUpsertBranchWithInsertOrUpdateFallbackReducesByKey() throws 
Exception {
+        RecordingDialect dialect = new RecordingDialect(false);
+
+        assertChangelogIsReducedByKey(dialect);
+
+        // no native upsert (Derby, Trino): exists + insert + update replace it
+        assertThat(dialect.upsertCalls).isEqualTo(1);
+        assertThat(dialect.rowExistsCalls).isEqualTo(1);
+        assertThat(dialect.insertCalls).isEqualTo(1);
+        assertThat(dialect.updateCalls).isEqualTo(1);
+        assertThat(dialect.deleteCalls).isEqualTo(1);
+    }
+
+    /**
+     * Shared by both upsert branches: within one buffer the last change to a 
key wins, the key
+     * alone decides identity, and a delete of an unknown key is a no-op.
+     */
+    private void assertChangelogIsReducedByKey(JdbcDialect dialect) throws 
Exception {
+        openOutputFormat(dialect, new String[] {"id"}, batchOf(100, 0), false);
+        TestEntry first = TEST_DATA[0];
+        TestEntry second = TEST_DATA[1];
+
+        outputFormat.writeRecord(changelogRow(RowKind.INSERT, first, "v1"));
+        outputFormat.writeRecord(changelogRow(RowKind.UPDATE_AFTER, first, 
"v2"));
+        outputFormat.flush();
+        assertThat(titlesById()).containsOnly(entry(first.id, "v2"));
+
+        // DELETE then INSERT of the same key keeps the row: the reduce key 
carries no row kind
+        outputFormat.writeRecord(changelogRow(RowKind.DELETE, first, "v2"));
+        outputFormat.writeRecord(changelogRow(RowKind.INSERT, first, "v3"));
+        outputFormat.flush();
+        assertThat(titlesById()).containsOnly(entry(first.id, "v3"));
+
+        // INSERT then DELETE of a new key writes nothing; a DELETE of an 
unknown key is a no-op
+        outputFormat.writeRecord(changelogRow(RowKind.INSERT, second, "v1"));
+        outputFormat.writeRecord(changelogRow(RowKind.DELETE, second, "v1"));
+        outputFormat.writeRecord(changelogRow(RowKind.DELETE, TEST_DATA[3], 
"never written"));
+        outputFormat.flush();
+        assertThat(titlesById()).containsOnly(entry(first.id, "v3"));
+
+        // an UPDATE_BEFORE on its own is a delete
+        outputFormat.writeRecord(changelogRow(RowKind.UPDATE_BEFORE, first, 
"v3"));
+        outputFormat.flush();
+        assertThat(titlesById()).isEmpty();
+    }
+
+    @Test
+    void testAppendOnlyBranchUsesThePlainInsert() throws Exception {
+        RecordingDialect dialect = new RecordingDialect(true);
+        openOutputFormat(dialect, null, batchOf(100, 0), false);
+
+        outputFormat.writeRecord(changelogRow(RowKind.INSERT, TEST_DATA[0], 
"a"));
+        outputFormat.writeRecord(changelogRow(RowKind.INSERT, TEST_DATA[1], 
"b"));
+        outputFormat.flush();
+
+        assertThat(titlesById())
+                .containsOnly(entry(TEST_DATA[0].id, "a"), 
entry(TEST_DATA[1].id, "b"));
+        assertThat(dialect.insertCalls).isEqualTo(1);
+        assertThat(dialect.upsertCalls).isZero();
+        assertThat(dialect.rowExistsCalls).isZero();
+        assertThat(dialect.deleteCalls).isZero();
+    }
+
+    @Test
+    void testKeyFieldOutsideTheFieldNamesFailsAtOpen() {
+        // indexOf gives -1 for the unknown key, and the builder indexes the 
field types with it
+        assertThatThrownBy(
+                        () ->
+                                openOutputFormat(
+                                        new DerbyDialect(),
+                                        new String[] {"nope"},
+                                        batchOf(100, 0),
+                                        false))
+                .isInstanceOf(ArrayIndexOutOfBoundsException.class);
+    }
+
+    @Test
+    void testNullKeyValueMatchesNeitherExistsNorDelete() throws Exception {
+        openOutputFormat(new DerbyDialect(), new String[] {"title"}, 
batchOf(100, 0), false);
+        TestEntry entry = TEST_DATA[0];
+
+        outputFormat.writeRecord(changelogRow(RowKind.INSERT, entry, null));
+        outputFormat.flush();
+        assertThat(titlesById()).containsOnly(entry(entry.id, null));
+
+        // `DELETE ... WHERE title = ?` with NULL matches nothing: the row is 
silently retained
+        outputFormat.writeRecord(changelogRow(RowKind.DELETE, entry, null));
+        outputFormat.flush();
+        assertThat(titlesById()).containsOnly(entry(entry.id, null));
+
+        // `exists` never matches either, so the same key is inserted again, 
which the table's
+        // primary key rejects
+        outputFormat.writeRecord(changelogRow(RowKind.INSERT, entry, null));
+        assertThatThrownBy(outputFormat::flush)
+                .isInstanceOf(IOException.class)
+                .hasCauseInstanceOf(SQLException.class);
+
+        // the buffer survived the failed flush: once the conflict is gone the 
replay lands
+        executeUpdate("DELETE FROM " + OUTPUT_TABLE_3 + " WHERE id = " + 
entry.id);
+        outputFormat.flush();
+        assertThat(titlesById()).containsOnly(entry(entry.id, null));
+    }
+
+    @Test
+    void testObjectReuseCopiesTheRecordBeforeBuffering() throws Exception {
+        openOutputFormat(new DerbyDialect(), new String[] {"id"}, batchOf(100, 
0), true);
+
+        GenericRowData row = (GenericRowData) changelogRow(RowKind.INSERT, 
TEST_DATA[0], "first");
+        outputFormat.writeRecord(row);
+        // with object reuse on, the runtime hands the same object over again 
for the next record
+        row.setField(0, TEST_DATA[1].id);
+        row.setField(1, StringData.fromString("second"));
+        outputFormat.writeRecord(row);
+        outputFormat.flush();
+
+        assertThat(titlesById())
+                .containsOnly(entry(TEST_DATA[0].id, "first"), 
entry(TEST_DATA[1].id, "second"));
+    }
+
+    @Test
+    void testFailedFlushKeepsTheBufferAndReplaysUpsertsIdempotently() throws 
Exception {

Review Comment:
   I've added `testFlushOnADeadConnectionReconnectsAndReplaysTheBuffer`, which 
closes the connection before the flush. With the `reconnect` flag ignored in 
`updateExecutor` it fails, and the FK test still passes.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to