eskabetxe commented on code in PR #246:
URL:
https://github.com/apache/flink-connector-jdbc/pull/246#discussion_r4135049564
##########
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:
Nice pin of the retry/replay behavior. One small gap: the FK conflict keeps
the connection valid, so this only drives updateExecutor(false); the
reconnect=true branch with the reduce executor (reestablish + re-prepare of
both statements) is covered only transitively. Closing the output format's
connection before the failing flush would pin that path exactly. Not a blocker.
--
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]