From 0e14729b66853dc6f7f4b9eb19b1b7156ba0aa8e Mon Sep 17 00:00:00 2001 From: Yonatan T Date: Wed, 2 Sep 2026 13:07:40 +0000 Subject: [PATCH 01/10] Support DELETE replication for tables without Primary Keys using ROWID - Fix regression in DatastreamToDML where tables lacking primary keys produced empty WHERE clauses for DELETE events. - Add ROWID fallback in getPrimaryKeyToValueFilterSql when primaryKeys is empty or unmapped. - Allow getDmlTemplate to generate DELETE DML when rowid or _metadata_row_id is present. - Make DatastreamRow getPrimaryKeys and getStringValue null-safe against missing metadata fields. - Add unit tests verifying ROWID fallback in DatastreamToDMLTest and DatastreamRowTest. Fixes b/506993545 --- .../v2/datastream/values/DatastreamRow.java | 13 ++- .../datastream/values/DatastreamRowTest.java | 25 ++++++ .../teleport/v2/utils/DatastreamToDML.java | 56 ++++++++++++- .../v2/utils/DatastreamToDMLTest.java | 81 +++++++++++++++++++ 4 files changed, 169 insertions(+), 6 deletions(-) diff --git a/v2/datastream-common/src/main/java/com/google/cloud/teleport/v2/datastream/values/DatastreamRow.java b/v2/datastream-common/src/main/java/com/google/cloud/teleport/v2/datastream/values/DatastreamRow.java index 29fe794cce..022812706f 100644 --- a/v2/datastream-common/src/main/java/com/google/cloud/teleport/v2/datastream/values/DatastreamRow.java +++ b/v2/datastream-common/src/main/java/com/google/cloud/teleport/v2/datastream/values/DatastreamRow.java @@ -90,7 +90,8 @@ public String getTableName() { public String getStringValue(String field) { if (this.jsonRow != null) { - return jsonRow.get(field).textValue(); + JsonNode node = jsonRow.get(field); + return node != null ? node.textValue() : null; } else { return (String) tableRow.get(field); } @@ -108,8 +109,12 @@ private Object getFieldValue(String field) { public List getPrimaryKeys() { List primaryKeys = new ArrayList<>(); if (this.jsonRow != null) { - for (JsonNode node : jsonRow.get("_metadata_primary_keys")) { - primaryKeys.add(node.asText()); + if (jsonRow.has("_metadata_primary_keys") + && jsonRow.get("_metadata_primary_keys") != null + && jsonRow.get("_metadata_primary_keys").isArray()) { + for (JsonNode node : jsonRow.get("_metadata_primary_keys")) { + primaryKeys.add(node.asText()); + } } } else { if (tableRow.get("_metadata_primary_keys") != null) { @@ -137,7 +142,7 @@ public List getPrimaryKeys() { } } - if (this.getSourceType().equals("oracle") && primaryKeys.isEmpty()) { + if ("oracle".equalsIgnoreCase(this.getSourceType()) && primaryKeys.isEmpty()) { primaryKeys.add(DEFAULT_ORACLE_PRIMARY_KEY); } diff --git a/v2/datastream-common/src/test/java/com/google/cloud/teleport/v2/datastream/values/DatastreamRowTest.java b/v2/datastream-common/src/test/java/com/google/cloud/teleport/v2/datastream/values/DatastreamRowTest.java index 28ab6831f9..8792f52257 100644 --- a/v2/datastream-common/src/test/java/com/google/cloud/teleport/v2/datastream/values/DatastreamRowTest.java +++ b/v2/datastream-common/src/test/java/com/google/cloud/teleport/v2/datastream/values/DatastreamRowTest.java @@ -16,7 +16,11 @@ package com.google.cloud.teleport.v2.datastream.values; import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertNull; +import static org.junit.Assert.assertTrue; +import com.fasterxml.jackson.databind.JsonNode; +import com.fasterxml.jackson.databind.ObjectMapper; import com.google.api.services.bigquery.model.TableRow; import java.io.IOException; import java.security.GeneralSecurityException; @@ -77,4 +81,25 @@ public void testSqlServerSortFields() { assertEquals("_metadata_timestamp", sortFields.get(0)); assertEquals("_metadata_lsn", sortFields.get(1)); } + + @Test + public void testGetPrimaryKeysJsonNode_nullSafeWhenNoMetadataField() throws IOException { + JsonNode jsonNode = new ObjectMapper().readTree("{\"id\": 123}"); + DatastreamRow row = DatastreamRow.of(jsonNode); + List pks = row.getPrimaryKeys(); + + assertTrue(pks.isEmpty()); + assertNull(row.getSourceType()); + } + + @Test + public void testGetPrimaryKeysJsonNode_oracleFallback() throws IOException { + JsonNode jsonNode = new ObjectMapper().readTree("{\"_metadata_source_type\": \"oracle\"}"); + DatastreamRow row = DatastreamRow.of(jsonNode); + List pks = row.getPrimaryKeys(); + + assertEquals(1, pks.size()); + assertEquals(DatastreamRow.DEFAULT_ORACLE_PRIMARY_KEY, pks.get(0)); + } } + diff --git a/v2/datastream-to-sql/src/main/java/com/google/cloud/teleport/v2/utils/DatastreamToDML.java b/v2/datastream-to-sql/src/main/java/com/google/cloud/teleport/v2/utils/DatastreamToDML.java index 5c5fbcad08..0e20e0dc0e 100644 --- a/v2/datastream-to-sql/src/main/java/com/google/cloud/teleport/v2/utils/DatastreamToDML.java +++ b/v2/datastream-to-sql/src/main/java/com/google/cloud/teleport/v2/utils/DatastreamToDML.java @@ -320,10 +320,20 @@ public DmlInfo convertJsonToDmlInfo(JsonNode rowObj, String failsafeValue) { failsafeValue); } + private boolean hasRowId(JsonNode rowObj) { + String casedRowId = applyCasingLogic(this.rowIdColumnName, this.columnCasing); + return rowObj.has(this.rowIdColumnName) + || rowObj.has(casedRowId) + || rowObj.has("_metadata_row_id"); + } + public String getDmlTemplate(JsonNode rowObj, List primaryKeys) { Boolean isDelete = rowObj.get("_metadata_deleted").asBoolean(); Boolean hasPrimaryKeys = primaryKeys.size() != 0; if (isDelete && !hasPrimaryKeys) { + if (hasRowId(rowObj)) { + return getDeleteDmlStatement(); + } throw new DeletedWithoutPrimaryKey("Delete DML without primary keys cannot be applied"); } else if (isDelete) { return getDeleteDmlStatement(); @@ -365,8 +375,14 @@ public String getValueSql(JsonNode rowObj, String columnName, Map it = rowObj.fieldNames(); it.hasNext(); ) { + String sourceFieldName = it.next(); + String destinationPkName = applyCasingLogic(sourceFieldName, this.columnCasing); + + if (primaryKeys.contains(destinationPkName)) { + String columnValue = getValueSql(rowObj, sourceFieldName, tableSchema); + String quotedDestinationPkName = quote(destinationPkName); + + if (pkToValueSql.isEmpty()) { + pkToValueSql = quotedDestinationPkName + "=" + columnValue; + } else { + pkToValueSql = pkToValueSql + " AND " + quotedDestinationPkName + "=" + columnValue; + } + } + } + } + + if (pkToValueSql.isEmpty() && hasRowId(rowObj)) { + String casedRowId = applyCasingLogic(this.rowIdColumnName, this.columnCasing); + String sourceRowIdField = null; + if (rowObj.has(this.rowIdColumnName)) { + sourceRowIdField = this.rowIdColumnName; + } else if (rowObj.has(casedRowId)) { + sourceRowIdField = casedRowId; + } else if (rowObj.has("_metadata_row_id")) { + sourceRowIdField = "_metadata_row_id"; + } + + if (sourceRowIdField != null) { + String columnValue = getValueSql(rowObj, sourceRowIdField, tableSchema); + pkToValueSql = quote(casedRowId) + "=" + columnValue; + } + } + return pkToValueSql; } diff --git a/v2/datastream-to-sql/src/test/java/com/google/cloud/teleport/v2/utils/DatastreamToDMLTest.java b/v2/datastream-to-sql/src/test/java/com/google/cloud/teleport/v2/utils/DatastreamToDMLTest.java index 6d2c75536a..bad7ba1b2a 100644 --- a/v2/datastream-to-sql/src/test/java/com/google/cloud/teleport/v2/utils/DatastreamToDMLTest.java +++ b/v2/datastream-to-sql/src/test/java/com/google/cloud/teleport/v2/utils/DatastreamToDMLTest.java @@ -1058,6 +1058,87 @@ public void testDmlTemplate_throwsErrorForDeleteWithoutPK() { assertThrows(RuntimeException.class, () -> dml.getDmlTemplate(rowObj, primaryKeys)); } + /** + * Tests that {@link DatastreamToDML#getDmlTemplate} returns a DELETE statement when + * primaryKeys is empty but the record contains a rowid. + */ + @Test + public void testDmlTemplate_supportsDeleteWithRowIdFallbackWhenNoPK() { + DatastreamToPostgresDML dml = DatastreamToPostgresDML.of(null); + String json = "{\"_metadata_deleted\": true, \"rowid\": \"AAAEARAAEAAAAC9AAS\"}"; + JsonNode rowObj = getRowObj(json); + List primaryKeys = java.util.Collections.emptyList(); + + String template = dml.getDmlTemplate(rowObj, primaryKeys); + assertEquals(dml.getDeleteDmlStatement(), template); + } + + /** + * Tests that {@link DatastreamToDML#getDmlTemplate} returns a DELETE statement when + * primaryKeys is empty but the record contains _metadata_row_id. + */ + @Test + public void testDmlTemplate_supportsDeleteWithMetadataRowIdFallbackWhenNoPK() { + DatastreamToPostgresDML dml = DatastreamToPostgresDML.of(null); + String json = "{\"_metadata_deleted\": true, \"_metadata_row_id\": \"AAAEARAAEAAAAC9AAS\"}"; + JsonNode rowObj = getRowObj(json); + List primaryKeys = java.util.Collections.emptyList(); + + String template = dml.getDmlTemplate(rowObj, primaryKeys); + assertEquals(dml.getDeleteDmlStatement(), template); + } + + /** + * Tests that {@link DatastreamToDML#getPrimaryKeyToValueFilterSql} generates the correct + * WHERE filter using rowid when primaryKeys contains rowid. + */ + @Test + public void testGetPrimaryKeyToValueFilterSql_usesRowIdWhenInPrimaryKeys() { + DatastreamToPostgresDML dml = DatastreamToPostgresDML.of(null); + String json = "{\"_metadata_deleted\": true, \"rowid\": \"AAAEARAAEAAAAC9AAS\"}"; + JsonNode rowObj = getRowObj(json); + List primaryKeys = Arrays.asList("rowid"); + Map tableSchema = new HashMap<>(); + tableSchema.put("rowid", "VARCHAR"); + + String filterSql = dml.getPrimaryKeyToValueFilterSql(rowObj, primaryKeys, tableSchema); + assertEquals("\"rowid\"='AAAEARAAEAAAAC9AAS'", filterSql); + } + + /** + * Tests that {@link DatastreamToDML#getPrimaryKeyToValueFilterSql} falls back to _metadata_row_id + * when primaryKeys has rowid and rowObj has _metadata_row_id. + */ + @Test + public void testGetPrimaryKeyToValueFilterSql_resolvesMetadataRowId() { + DatastreamToPostgresDML dml = DatastreamToPostgresDML.of(null); + String json = "{\"_metadata_deleted\": true, \"_metadata_row_id\": \"AAAEARAAEAAAAC9AAS\"}"; + JsonNode rowObj = getRowObj(json); + List primaryKeys = Arrays.asList("rowid"); + Map tableSchema = new HashMap<>(); + tableSchema.put("rowid", "VARCHAR"); + + String filterSql = dml.getPrimaryKeyToValueFilterSql(rowObj, primaryKeys, tableSchema); + assertEquals("\"rowid\"='AAAEARAAEAAAAC9AAS'", filterSql); + } + + /** + * Tests that {@link DatastreamToDML#getPrimaryKeyToValueFilterSql} generates a rowid filter + * even when destination primaryKeys list is empty (table without PK). + */ + @Test + public void testGetPrimaryKeyToValueFilterSql_usesRowIdFallbackWhenNoPK() { + DatastreamToPostgresDML dml = DatastreamToPostgresDML.of(null); + String json = "{\"_metadata_deleted\": true, \"rowid\": \"AAAEARAAEAAAAC9AAS\"}"; + JsonNode rowObj = getRowObj(json); + List primaryKeys = java.util.Collections.emptyList(); + Map tableSchema = new HashMap<>(); + tableSchema.put("rowid", "VARCHAR"); + + String filterSql = dml.getPrimaryKeyToValueFilterSql(rowObj, primaryKeys, tableSchema); + assertEquals("\"rowid\"='AAAEARAAEAAAAC9AAS'", filterSql); + } + /** * Tests basic numeric type cleansing in Postgres, specifically ensuring empty strings become * NULL. From cc17e7f06ff75fff11ee95a60c356352ec0bb7c4 Mon Sep 17 00:00:00 2001 From: Yonatan T Date: Tue, 8 Sep 2026 12:56:08 +0000 Subject: [PATCH 02/10] Apply spotless formatting to DatastreamRowTest and DatastreamToDMLTest --- .../v2/datastream/values/DatastreamRowTest.java | 1 - .../teleport/v2/utils/DatastreamToDMLTest.java | 16 ++++++++-------- 2 files changed, 8 insertions(+), 9 deletions(-) diff --git a/v2/datastream-common/src/test/java/com/google/cloud/teleport/v2/datastream/values/DatastreamRowTest.java b/v2/datastream-common/src/test/java/com/google/cloud/teleport/v2/datastream/values/DatastreamRowTest.java index 8792f52257..ea7ed30d0e 100644 --- a/v2/datastream-common/src/test/java/com/google/cloud/teleport/v2/datastream/values/DatastreamRowTest.java +++ b/v2/datastream-common/src/test/java/com/google/cloud/teleport/v2/datastream/values/DatastreamRowTest.java @@ -102,4 +102,3 @@ public void testGetPrimaryKeysJsonNode_oracleFallback() throws IOException { assertEquals(DatastreamRow.DEFAULT_ORACLE_PRIMARY_KEY, pks.get(0)); } } - diff --git a/v2/datastream-to-sql/src/test/java/com/google/cloud/teleport/v2/utils/DatastreamToDMLTest.java b/v2/datastream-to-sql/src/test/java/com/google/cloud/teleport/v2/utils/DatastreamToDMLTest.java index bad7ba1b2a..f3194e3786 100644 --- a/v2/datastream-to-sql/src/test/java/com/google/cloud/teleport/v2/utils/DatastreamToDMLTest.java +++ b/v2/datastream-to-sql/src/test/java/com/google/cloud/teleport/v2/utils/DatastreamToDMLTest.java @@ -1059,8 +1059,8 @@ public void testDmlTemplate_throwsErrorForDeleteWithoutPK() { } /** - * Tests that {@link DatastreamToDML#getDmlTemplate} returns a DELETE statement when - * primaryKeys is empty but the record contains a rowid. + * Tests that {@link DatastreamToDML#getDmlTemplate} returns a DELETE statement when primaryKeys + * is empty but the record contains a rowid. */ @Test public void testDmlTemplate_supportsDeleteWithRowIdFallbackWhenNoPK() { @@ -1074,8 +1074,8 @@ public void testDmlTemplate_supportsDeleteWithRowIdFallbackWhenNoPK() { } /** - * Tests that {@link DatastreamToDML#getDmlTemplate} returns a DELETE statement when - * primaryKeys is empty but the record contains _metadata_row_id. + * Tests that {@link DatastreamToDML#getDmlTemplate} returns a DELETE statement when primaryKeys + * is empty but the record contains _metadata_row_id. */ @Test public void testDmlTemplate_supportsDeleteWithMetadataRowIdFallbackWhenNoPK() { @@ -1089,8 +1089,8 @@ public void testDmlTemplate_supportsDeleteWithMetadataRowIdFallbackWhenNoPK() { } /** - * Tests that {@link DatastreamToDML#getPrimaryKeyToValueFilterSql} generates the correct - * WHERE filter using rowid when primaryKeys contains rowid. + * Tests that {@link DatastreamToDML#getPrimaryKeyToValueFilterSql} generates the correct WHERE + * filter using rowid when primaryKeys contains rowid. */ @Test public void testGetPrimaryKeyToValueFilterSql_usesRowIdWhenInPrimaryKeys() { @@ -1123,8 +1123,8 @@ public void testGetPrimaryKeyToValueFilterSql_resolvesMetadataRowId() { } /** - * Tests that {@link DatastreamToDML#getPrimaryKeyToValueFilterSql} generates a rowid filter - * even when destination primaryKeys list is empty (table without PK). + * Tests that {@link DatastreamToDML#getPrimaryKeyToValueFilterSql} generates a rowid filter even + * when destination primaryKeys list is empty (table without PK). */ @Test public void testGetPrimaryKeyToValueFilterSql_usesRowIdFallbackWhenNoPK() { From 13be64a5a14fc3542d161ad8c6fe8f0df9730d4a Mon Sep 17 00:00:00 2001 From: Yonatan T Date: Tue, 8 Sep 2026 13:32:05 +0000 Subject: [PATCH 03/10] Address review comments: handle column casing in row ID fallback and optimize primary key lookup --- .../v2/datastream/values/DatastreamRow.java | 7 ++-- .../teleport/v2/utils/DatastreamToDML.java | 12 +++++-- .../v2/utils/DatastreamToDMLTest.java | 35 +++++++++++++++++++ 3 files changed, 47 insertions(+), 7 deletions(-) diff --git a/v2/datastream-common/src/main/java/com/google/cloud/teleport/v2/datastream/values/DatastreamRow.java b/v2/datastream-common/src/main/java/com/google/cloud/teleport/v2/datastream/values/DatastreamRow.java index 022812706f..62e3987697 100644 --- a/v2/datastream-common/src/main/java/com/google/cloud/teleport/v2/datastream/values/DatastreamRow.java +++ b/v2/datastream-common/src/main/java/com/google/cloud/teleport/v2/datastream/values/DatastreamRow.java @@ -109,10 +109,9 @@ private Object getFieldValue(String field) { public List getPrimaryKeys() { List primaryKeys = new ArrayList<>(); if (this.jsonRow != null) { - if (jsonRow.has("_metadata_primary_keys") - && jsonRow.get("_metadata_primary_keys") != null - && jsonRow.get("_metadata_primary_keys").isArray()) { - for (JsonNode node : jsonRow.get("_metadata_primary_keys")) { + JsonNode pksNode = jsonRow.get("_metadata_primary_keys"); + if (pksNode != null && pksNode.isArray()) { + for (JsonNode node : pksNode) { primaryKeys.add(node.asText()); } } diff --git a/v2/datastream-to-sql/src/main/java/com/google/cloud/teleport/v2/utils/DatastreamToDML.java b/v2/datastream-to-sql/src/main/java/com/google/cloud/teleport/v2/utils/DatastreamToDML.java index 0e20e0dc0e..75ae26dfae 100644 --- a/v2/datastream-to-sql/src/main/java/com/google/cloud/teleport/v2/utils/DatastreamToDML.java +++ b/v2/datastream-to-sql/src/main/java/com/google/cloud/teleport/v2/utils/DatastreamToDML.java @@ -375,10 +375,16 @@ public String getValueSql(JsonNode rowObj, String columnName, Map primaryKeys = java.util.Collections.emptyList(); + Map tableSchema = new HashMap<>(); + tableSchema.put("ROWID", "VARCHAR"); + + String filterSql = dml.getPrimaryKeyToValueFilterSql(rowObj, primaryKeys, tableSchema); + assertEquals("\"ROWID\"='AAAEARAAEAAAAC9AAS'", filterSql); + } + + /** + * Tests that {@link DatastreamToDML#getValueSql} resolves cased rowid when column is UPPERCASE + * and record contains _metadata_row_id. + */ + @Test + public void testGetValueSql_supportsCasedRowIdFallback() { + DatastreamToPostgresDML dml = + (DatastreamToPostgresDML) DatastreamToPostgresDML.of(null).withColumnCasing("UPPERCASE"); + String json = "{\"_metadata_row_id\": \"AAAEARAAEAAAAC9AAS\"}"; + JsonNode rowObj = getRowObj(json); + Map tableSchema = new HashMap<>(); + tableSchema.put("ROWID", "VARCHAR"); + + String valueSql = dml.getValueSql(rowObj, "ROWID", tableSchema); + assertEquals("'AAAEARAAEAAAAC9AAS'", valueSql); + } + /** * Tests basic numeric type cleansing in Postgres, specifically ensuring empty strings become * NULL. From 2ffbde8cc4e2d7910924510f33e9fc69491ae67a Mon Sep 17 00:00:00 2001 From: Yonatan T Date: Wed, 9 Sep 2026 05:11:18 +0000 Subject: [PATCH 04/10] test: add comprehensive unit tests for row ID fallback and casing to achieve 100% patch coverage --- .../datastream/values/DatastreamRowTest.java | 23 ++++- .../v2/utils/DatastreamToDMLTest.java | 85 +++++++++++++++++++ 2 files changed, 106 insertions(+), 2 deletions(-) diff --git a/v2/datastream-common/src/test/java/com/google/cloud/teleport/v2/datastream/values/DatastreamRowTest.java b/v2/datastream-common/src/test/java/com/google/cloud/teleport/v2/datastream/values/DatastreamRowTest.java index ea7ed30d0e..0d2e2bba33 100644 --- a/v2/datastream-common/src/test/java/com/google/cloud/teleport/v2/datastream/values/DatastreamRowTest.java +++ b/v2/datastream-common/src/test/java/com/google/cloud/teleport/v2/datastream/values/DatastreamRowTest.java @@ -93,11 +93,30 @@ public void testGetPrimaryKeysJsonNode_nullSafeWhenNoMetadataField() throws IOEx } @Test - public void testGetPrimaryKeysJsonNode_oracleFallback() throws IOException { - JsonNode jsonNode = new ObjectMapper().readTree("{\"_metadata_source_type\": \"oracle\"}"); + public void testGetPrimaryKeysJsonNode_withPrimaryKeysArray() throws IOException { + JsonNode jsonNode = + new ObjectMapper() + .readTree( + "{\"_metadata_source_type\": \"oracle\", \"_metadata_primary_keys\":" + + " [\"id\", \"name\"]}"); DatastreamRow row = DatastreamRow.of(jsonNode); List pks = row.getPrimaryKeys(); + assertEquals(2, pks.size()); + assertEquals("id", pks.get(0)); + assertEquals("name", pks.get(1)); + } + + @Test + public void testGetPrimaryKeysJsonNode_nonArrayPrimaryKeysNode() throws IOException { + JsonNode jsonNode = + new ObjectMapper() + .readTree( + "{\"_metadata_source_type\": \"oracle\", \"_metadata_primary_keys\": \"id\"}"); + DatastreamRow row = DatastreamRow.of(jsonNode); + List pks = row.getPrimaryKeys(); + + // Non-array node should not be treated as array, falls back to default oracle PK assertEquals(1, pks.size()); assertEquals(DatastreamRow.DEFAULT_ORACLE_PRIMARY_KEY, pks.get(0)); } diff --git a/v2/datastream-to-sql/src/test/java/com/google/cloud/teleport/v2/utils/DatastreamToDMLTest.java b/v2/datastream-to-sql/src/test/java/com/google/cloud/teleport/v2/utils/DatastreamToDMLTest.java index a8b6367ee5..770615b3e4 100644 --- a/v2/datastream-to-sql/src/test/java/com/google/cloud/teleport/v2/utils/DatastreamToDMLTest.java +++ b/v2/datastream-to-sql/src/test/java/com/google/cloud/teleport/v2/utils/DatastreamToDMLTest.java @@ -1174,6 +1174,91 @@ public void testGetValueSql_supportsCasedRowIdFallback() { assertEquals("'AAAEARAAEAAAAC9AAS'", valueSql); } + /** + * Tests that {@link DatastreamToDML#getValueSql} resolves _metadata_row_id when column is + * _metadata_row_id and record contains rowid. + */ + @Test + public void testGetValueSql_supportsMetadataRowIdFallbackFromRowId() { + DatastreamToPostgresDML dml = DatastreamToPostgresDML.of(null); + String json = "{\"rowid\": \"AAAEARAAEAAAAC9AAS\"}"; + JsonNode rowObj = getRowObj(json); + Map tableSchema = new HashMap<>(); + tableSchema.put("_metadata_row_id", "VARCHAR"); + + String valueSql = dml.getValueSql(rowObj, "_metadata_row_id", tableSchema); + assertEquals("'AAAEARAAEAAAAC9AAS'", valueSql); + } + + /** + * Tests that {@link DatastreamToDML#getValueSql} resolves _metadata_row_id when column is + * _metadata_row_id and record contains cased ROWID. + */ + @Test + public void testGetValueSql_supportsMetadataRowIdFallbackFromCasedRowId() { + DatastreamToPostgresDML dml = + (DatastreamToPostgresDML) DatastreamToPostgresDML.of(null).withColumnCasing("UPPERCASE"); + String json = "{\"ROWID\": \"AAAEARAAEAAAAC9AAS\"}"; + JsonNode rowObj = getRowObj(json); + Map tableSchema = new HashMap<>(); + tableSchema.put("_metadata_row_id", "VARCHAR"); + + String valueSql = dml.getValueSql(rowObj, "_metadata_row_id", tableSchema); + assertEquals("'AAAEARAAEAAAAC9AAS'", valueSql); + } + + /** + * Tests that {@link DatastreamToDML#getValueSql} warns and returns empty string when requested + * column is missing and not a rowid. + */ + @Test + public void testGetValueSql_returnsEmptyOnMissingColumn() { + DatastreamToPostgresDML dml = DatastreamToPostgresDML.of(null); + String json = "{\"id\": 123}"; + JsonNode rowObj = getRowObj(json); + Map tableSchema = new HashMap<>(); + + String valueSql = dml.getValueSql(rowObj, "missing_column", tableSchema); + assertEquals("", valueSql); + } + + /** + * Tests that {@link DatastreamToDML#getPrimaryKeyToValueFilterSql} generates a cased rowid filter + * when destination primaryKeys is empty and record contains cased ROWID. + */ + @Test + public void testGetPrimaryKeyToValueFilterSql_usesCasedRowIdInRecordWhenNoPK() { + DatastreamToPostgresDML dml = + (DatastreamToPostgresDML) DatastreamToPostgresDML.of(null).withColumnCasing("UPPERCASE"); + String json = "{\"_metadata_deleted\": true, \"ROWID\": \"AAAEARAAEAAAAC9AAS\"}"; + JsonNode rowObj = getRowObj(json); + List primaryKeys = java.util.Collections.emptyList(); + Map tableSchema = new HashMap<>(); + tableSchema.put("ROWID", "VARCHAR"); + + String filterSql = dml.getPrimaryKeyToValueFilterSql(rowObj, primaryKeys, tableSchema); + assertEquals("\"ROWID\"='AAAEARAAEAAAAC9AAS'", filterSql); + } + + /** + * Tests that {@link DatastreamToDML#getPrimaryKeyToValueFilterSql} joins multiple cased primary + * keys with AND. + */ + @Test + public void testGetPrimaryKeyToValueFilterSql_multipleCasedPrimaryKeys() { + DatastreamToPostgresDML dml = + (DatastreamToPostgresDML) DatastreamToPostgresDML.of(null).withColumnCasing("UPPERCASE"); + String json = "{\"id\": 1, \"org_id\": 2}"; + JsonNode rowObj = getRowObj(json); + List primaryKeys = Arrays.asList("ID", "ORG_ID"); + Map tableSchema = new HashMap<>(); + tableSchema.put("ID", "INTEGER"); + tableSchema.put("ORG_ID", "INTEGER"); + + String filterSql = dml.getPrimaryKeyToValueFilterSql(rowObj, primaryKeys, tableSchema); + assertEquals("\"ID\"=1 AND \"ORG_ID\"=2", filterSql); + } + /** * Tests basic numeric type cleansing in Postgres, specifically ensuring empty strings become * NULL. From 72f2afbe51db22bafccd4d1661043498559c9a44 Mon Sep 17 00:00:00 2001 From: Yonatan T Date: Wed, 9 Sep 2026 07:23:06 +0000 Subject: [PATCH 05/10] Validate target schema contains rowid column before generating rowid fallback WHERE filter --- .../teleport/v2/utils/DatastreamToDML.java | 8 +++++++- .../v2/utils/DatastreamToDMLTest.java | 19 +++++++++++++++++++ 2 files changed, 26 insertions(+), 1 deletion(-) diff --git a/v2/datastream-to-sql/src/main/java/com/google/cloud/teleport/v2/utils/DatastreamToDML.java b/v2/datastream-to-sql/src/main/java/com/google/cloud/teleport/v2/utils/DatastreamToDML.java index 75ae26dfae..77ce4d9b32 100644 --- a/v2/datastream-to-sql/src/main/java/com/google/cloud/teleport/v2/utils/DatastreamToDML.java +++ b/v2/datastream-to-sql/src/main/java/com/google/cloud/teleport/v2/utils/DatastreamToDML.java @@ -81,7 +81,7 @@ public abstract class DatastreamToDML public abstract String getTargetSchemaName(DatastreamRow row); /* An exception for delete DML without a primary key */ - private class DeletedWithoutPrimaryKey extends RuntimeException { + static class DeletedWithoutPrimaryKey extends RuntimeException { public DeletedWithoutPrimaryKey(String errorMessage) { super(errorMessage); } @@ -545,6 +545,12 @@ public String getPrimaryKeyToValueFilterSql( if (pkToValueSql.isEmpty() && hasRowId(rowObj)) { String casedRowId = applyCasingLogic(this.rowIdColumnName, this.columnCasing); + if (!tableSchema.containsKey(casedRowId)) { + throw new DeletedWithoutPrimaryKey( + String.format( + "Cannot replicate DELETE for table without primary keys: target schema lacks '%s' column.", + casedRowId)); + } String sourceRowIdField = null; if (rowObj.has(this.rowIdColumnName)) { sourceRowIdField = this.rowIdColumnName; diff --git a/v2/datastream-to-sql/src/test/java/com/google/cloud/teleport/v2/utils/DatastreamToDMLTest.java b/v2/datastream-to-sql/src/test/java/com/google/cloud/teleport/v2/utils/DatastreamToDMLTest.java index 770615b3e4..0ea858df4d 100644 --- a/v2/datastream-to-sql/src/test/java/com/google/cloud/teleport/v2/utils/DatastreamToDMLTest.java +++ b/v2/datastream-to-sql/src/test/java/com/google/cloud/teleport/v2/utils/DatastreamToDMLTest.java @@ -1139,6 +1139,25 @@ public void testGetPrimaryKeyToValueFilterSql_usesRowIdFallbackWhenNoPK() { assertEquals("\"rowid\"='AAAEARAAEAAAAC9AAS'", filterSql); } + /** + * Tests that {@link DatastreamToDML#getPrimaryKeyToValueFilterSql} throws {@link + * DatastreamToDML.DeletedWithoutPrimaryKey} when destination primaryKeys is empty and target + * schema lacks the rowid column. + */ + @Test + public void testGetPrimaryKeyToValueFilterSql_throwsWhenTargetSchemaLacksRowId() { + DatastreamToPostgresDML dml = DatastreamToPostgresDML.of(null); + String json = "{\"_metadata_deleted\": true, \"rowid\": \"AAAEARAAEAAAAC9AAS\"}"; + JsonNode rowObj = getRowObj(json); + List primaryKeys = java.util.Collections.emptyList(); + Map tableSchema = new HashMap<>(); + tableSchema.put("name", "VARCHAR"); + + assertThrows( + DatastreamToDML.DeletedWithoutPrimaryKey.class, + () -> dml.getPrimaryKeyToValueFilterSql(rowObj, primaryKeys, tableSchema)); + } + /** * Tests that {@link DatastreamToDML#getPrimaryKeyToValueFilterSql} generates a cased rowid filter * with UPPERCASE column casing when destination primaryKeys is empty. From e064af99d8f628f5ac4c97f96ea623cc484779cd Mon Sep 17 00:00:00 2001 From: Yonatan T Date: Tue, 22 Sep 2026 09:32:53 +0000 Subject: [PATCH 06/10] Revert changes in datastream-common to keep PR scoped to datastream-to-sql --- .../v2/datastream/values/DatastreamRow.java | 12 ++---- .../datastream/values/DatastreamRowTest.java | 43 ------------------- 2 files changed, 4 insertions(+), 51 deletions(-) diff --git a/v2/datastream-common/src/main/java/com/google/cloud/teleport/v2/datastream/values/DatastreamRow.java b/v2/datastream-common/src/main/java/com/google/cloud/teleport/v2/datastream/values/DatastreamRow.java index 62e3987697..29fe794cce 100644 --- a/v2/datastream-common/src/main/java/com/google/cloud/teleport/v2/datastream/values/DatastreamRow.java +++ b/v2/datastream-common/src/main/java/com/google/cloud/teleport/v2/datastream/values/DatastreamRow.java @@ -90,8 +90,7 @@ public String getTableName() { public String getStringValue(String field) { if (this.jsonRow != null) { - JsonNode node = jsonRow.get(field); - return node != null ? node.textValue() : null; + return jsonRow.get(field).textValue(); } else { return (String) tableRow.get(field); } @@ -109,11 +108,8 @@ private Object getFieldValue(String field) { public List getPrimaryKeys() { List primaryKeys = new ArrayList<>(); if (this.jsonRow != null) { - JsonNode pksNode = jsonRow.get("_metadata_primary_keys"); - if (pksNode != null && pksNode.isArray()) { - for (JsonNode node : pksNode) { - primaryKeys.add(node.asText()); - } + for (JsonNode node : jsonRow.get("_metadata_primary_keys")) { + primaryKeys.add(node.asText()); } } else { if (tableRow.get("_metadata_primary_keys") != null) { @@ -141,7 +137,7 @@ public List getPrimaryKeys() { } } - if ("oracle".equalsIgnoreCase(this.getSourceType()) && primaryKeys.isEmpty()) { + if (this.getSourceType().equals("oracle") && primaryKeys.isEmpty()) { primaryKeys.add(DEFAULT_ORACLE_PRIMARY_KEY); } diff --git a/v2/datastream-common/src/test/java/com/google/cloud/teleport/v2/datastream/values/DatastreamRowTest.java b/v2/datastream-common/src/test/java/com/google/cloud/teleport/v2/datastream/values/DatastreamRowTest.java index 0d2e2bba33..28ab6831f9 100644 --- a/v2/datastream-common/src/test/java/com/google/cloud/teleport/v2/datastream/values/DatastreamRowTest.java +++ b/v2/datastream-common/src/test/java/com/google/cloud/teleport/v2/datastream/values/DatastreamRowTest.java @@ -16,11 +16,7 @@ package com.google.cloud.teleport.v2.datastream.values; import static org.junit.Assert.assertEquals; -import static org.junit.Assert.assertNull; -import static org.junit.Assert.assertTrue; -import com.fasterxml.jackson.databind.JsonNode; -import com.fasterxml.jackson.databind.ObjectMapper; import com.google.api.services.bigquery.model.TableRow; import java.io.IOException; import java.security.GeneralSecurityException; @@ -81,43 +77,4 @@ public void testSqlServerSortFields() { assertEquals("_metadata_timestamp", sortFields.get(0)); assertEquals("_metadata_lsn", sortFields.get(1)); } - - @Test - public void testGetPrimaryKeysJsonNode_nullSafeWhenNoMetadataField() throws IOException { - JsonNode jsonNode = new ObjectMapper().readTree("{\"id\": 123}"); - DatastreamRow row = DatastreamRow.of(jsonNode); - List pks = row.getPrimaryKeys(); - - assertTrue(pks.isEmpty()); - assertNull(row.getSourceType()); - } - - @Test - public void testGetPrimaryKeysJsonNode_withPrimaryKeysArray() throws IOException { - JsonNode jsonNode = - new ObjectMapper() - .readTree( - "{\"_metadata_source_type\": \"oracle\", \"_metadata_primary_keys\":" - + " [\"id\", \"name\"]}"); - DatastreamRow row = DatastreamRow.of(jsonNode); - List pks = row.getPrimaryKeys(); - - assertEquals(2, pks.size()); - assertEquals("id", pks.get(0)); - assertEquals("name", pks.get(1)); - } - - @Test - public void testGetPrimaryKeysJsonNode_nonArrayPrimaryKeysNode() throws IOException { - JsonNode jsonNode = - new ObjectMapper() - .readTree( - "{\"_metadata_source_type\": \"oracle\", \"_metadata_primary_keys\": \"id\"}"); - DatastreamRow row = DatastreamRow.of(jsonNode); - List pks = row.getPrimaryKeys(); - - // Non-array node should not be treated as array, falls back to default oracle PK - assertEquals(1, pks.size()); - assertEquals(DatastreamRow.DEFAULT_ORACLE_PRIMARY_KEY, pks.get(0)); - } } From 8de4fed4d3f6f36fd3e749759a7be5db8aaa0dc2 Mon Sep 17 00:00:00 2001 From: Yonatan T Date: Tue, 22 Sep 2026 12:32:32 +0000 Subject: [PATCH 07/10] fix(datastream-to-sql): null-safe source primary key extraction in DatastreamToDML --- .../teleport/v2/utils/DatastreamToDML.java | 19 +++++++++++++++++-- 1 file changed, 17 insertions(+), 2 deletions(-) diff --git a/v2/datastream-to-sql/src/main/java/com/google/cloud/teleport/v2/utils/DatastreamToDML.java b/v2/datastream-to-sql/src/main/java/com/google/cloud/teleport/v2/utils/DatastreamToDML.java index 77ce4d9b32..116f9dc89e 100644 --- a/v2/datastream-to-sql/src/main/java/com/google/cloud/teleport/v2/utils/DatastreamToDML.java +++ b/v2/datastream-to-sql/src/main/java/com/google/cloud/teleport/v2/utils/DatastreamToDML.java @@ -30,6 +30,7 @@ import java.sql.SQLException; import java.util.ArrayList; import java.util.Arrays; +import java.util.Collections; import java.util.HashMap; import java.util.Iterator; import java.util.List; @@ -503,11 +504,25 @@ public String getColumnsUpdateSql(JsonNode rowObj, Map tableSche return onUpdateSql; } + private List getSourcePrimaryKeys(JsonNode rowObj) { + if (rowObj == null) { + return Collections.emptyList(); + } + JsonNode pksNode = rowObj.get("_metadata_primary_keys"); + if (pksNode != null && pksNode.isArray()) { + List primaryKeys = new ArrayList<>(); + for (JsonNode node : pksNode) { + primaryKeys.add(node.asText()); + } + return primaryKeys; + } + return Collections.emptyList(); + } + public String getPrimaryKeyToValueFilterSql( JsonNode rowObj, List primaryKeys, Map tableSchema) { - DatastreamRow row = DatastreamRow.of(rowObj); - List sourcePrimaryKeys = row.getPrimaryKeys(); + List sourcePrimaryKeys = getSourcePrimaryKeys(rowObj); String pkToValueSql = ""; for (String sourcePkName : sourcePrimaryKeys) { From 26d18a56f07cb51a83b170a979c17baf49d87a4e Mon Sep 17 00:00:00 2001 From: Yonatan T Date: Wed, 23 Sep 2026 12:49:40 +0000 Subject: [PATCH 08/10] refactor(datastream-to-sql): simplify rowid fallback and support destination rowid primary key --- .../teleport/v2/utils/DatastreamToDML.java | 61 ++++++++-------- .../v2/utils/DatastreamToDMLTest.java | 71 ++++++++++--------- 2 files changed, 70 insertions(+), 62 deletions(-) diff --git a/v2/datastream-to-sql/src/main/java/com/google/cloud/teleport/v2/utils/DatastreamToDML.java b/v2/datastream-to-sql/src/main/java/com/google/cloud/teleport/v2/utils/DatastreamToDML.java index 116f9dc89e..05a6a2cb5e 100644 --- a/v2/datastream-to-sql/src/main/java/com/google/cloud/teleport/v2/utils/DatastreamToDML.java +++ b/v2/datastream-to-sql/src/main/java/com/google/cloud/teleport/v2/utils/DatastreamToDML.java @@ -36,6 +36,7 @@ import java.util.List; import java.util.Map; import java.util.Objects; +import java.util.stream.Stream; import javax.sql.DataSource; import org.apache.beam.sdk.transforms.DoFn; import org.apache.beam.sdk.values.KV; @@ -270,6 +271,9 @@ public List getPrimaryKeys( for (String destPk : destinationPrimaryKeys) { if (!casedSourceFieldNames.contains(destPk)) { + if (destPk.equalsIgnoreCase(this.rowIdColumnName) && rowObj.has("_metadata_row_id")) { + continue; + } return this.getDefaultPrimaryKeys(); } } @@ -373,29 +377,19 @@ public Map getSqlTemplateValues( } public String getValueSql(JsonNode rowObj, String columnName, Map tableSchema) { - String columnValue; + if (columnName.equalsIgnoreCase(this.rowIdColumnName) + && !rowObj.has(columnName) + && rowObj.has("_metadata_row_id")) { + columnName = "_metadata_row_id"; + } + JsonNode columnObj = rowObj.get(columnName); if (columnObj == null) { - String casedRowId = applyCasingLogic(this.rowIdColumnName, this.columnCasing); - if ((columnName.equals(this.rowIdColumnName) || columnName.equals(casedRowId)) - && rowObj.has("_metadata_row_id")) { - columnObj = rowObj.get("_metadata_row_id"); - } else if (columnName.equals("_metadata_row_id") - && (rowObj.has(this.rowIdColumnName) || rowObj.has(casedRowId))) { - columnObj = - rowObj.has(this.rowIdColumnName) - ? rowObj.get(this.rowIdColumnName) - : rowObj.get(casedRowId); - } else { - LOG.warn("Missing Required Value: {} in {}", columnName, rowObj.toString()); - return ""; - } - } - if (columnObj.isTextual()) { - columnValue = "\'" + cleanSql(columnObj.textValue()) + "\'"; - } else { - columnValue = columnObj.toString(); + LOG.warn("Missing Required Value: {} in {}", columnName, rowObj.toString()); + return ""; } + String columnValue = + columnObj.isTextual() ? "'" + cleanSql(columnObj.textValue()) + "'" : columnObj.toString(); return cleanDataTypeValueSql(columnValue, columnName, tableSchema); } @@ -556,9 +550,21 @@ public String getPrimaryKeyToValueFilterSql( } } } + + if (pkToValueSql.isEmpty() && rowObj.has("_metadata_row_id")) { + for (String pk : primaryKeys) { + if (pk.equalsIgnoreCase(this.rowIdColumnName)) { + String columnValue = getValueSql(rowObj, pk, tableSchema); + pkToValueSql = quote(pk) + "=" + columnValue; + break; + } + } + } } - if (pkToValueSql.isEmpty() && hasRowId(rowObj)) { + boolean isDelete = + rowObj.has("_metadata_deleted") && rowObj.get("_metadata_deleted").asBoolean(); + if (isDelete && pkToValueSql.isEmpty() && primaryKeys.isEmpty()) { String casedRowId = applyCasingLogic(this.rowIdColumnName, this.columnCasing); if (!tableSchema.containsKey(casedRowId)) { throw new DeletedWithoutPrimaryKey( @@ -566,14 +572,11 @@ public String getPrimaryKeyToValueFilterSql( "Cannot replicate DELETE for table without primary keys: target schema lacks '%s' column.", casedRowId)); } - String sourceRowIdField = null; - if (rowObj.has(this.rowIdColumnName)) { - sourceRowIdField = this.rowIdColumnName; - } else if (rowObj.has(casedRowId)) { - sourceRowIdField = casedRowId; - } else if (rowObj.has("_metadata_row_id")) { - sourceRowIdField = "_metadata_row_id"; - } + String sourceRowIdField = + Stream.of(casedRowId, this.rowIdColumnName, "_metadata_row_id") + .filter(rowObj::has) + .findFirst() + .orElse(null); if (sourceRowIdField != null) { String columnValue = getValueSql(rowObj, sourceRowIdField, tableSchema); diff --git a/v2/datastream-to-sql/src/test/java/com/google/cloud/teleport/v2/utils/DatastreamToDMLTest.java b/v2/datastream-to-sql/src/test/java/com/google/cloud/teleport/v2/utils/DatastreamToDMLTest.java index 0ea858df4d..09ca38aee9 100644 --- a/v2/datastream-to-sql/src/test/java/com/google/cloud/teleport/v2/utils/DatastreamToDMLTest.java +++ b/v2/datastream-to-sql/src/test/java/com/google/cloud/teleport/v2/utils/DatastreamToDMLTest.java @@ -1122,6 +1122,44 @@ public void testGetPrimaryKeyToValueFilterSql_resolvesMetadataRowId() { assertEquals("\"rowid\"='AAAEARAAEAAAAC9AAS'", filterSql); } + /** + * Tests that {@link DatastreamToDML#getPrimaryKeyToValueFilterSql} resolves _metadata_row_id when + * primaryKeys has cased ROWID and rowObj has _metadata_row_id. + */ + @Test + public void testGetPrimaryKeyToValueFilterSql_resolvesMetadataRowIdWithCasedRowIdPk() { + DatastreamToPostgresDML dml = + (DatastreamToPostgresDML) DatastreamToPostgresDML.of(null).withColumnCasing("UPPERCASE"); + String json = "{\"_metadata_deleted\": true, \"_metadata_row_id\": \"AAAEARAAEAAAAC9AAS\"}"; + JsonNode rowObj = getRowObj(json); + List primaryKeys = Arrays.asList("ROWID"); + Map tableSchema = new HashMap<>(); + tableSchema.put("ROWID", "VARCHAR"); + + String filterSql = dml.getPrimaryKeyToValueFilterSql(rowObj, primaryKeys, tableSchema); + assertEquals("\"ROWID\"='AAAEARAAEAAAAC9AAS'", filterSql); + } + + /** + * Tests that {@link DatastreamToDML#getPrimaryKeyToValueFilterSql} resolves _metadata_row_id when + * _metadata_primary_keys contains ROWID and rowObj has _metadata_row_id. + */ + @Test + public void testGetPrimaryKeyToValueFilterSql_resolvesMetadataRowIdWithSourcePKs() { + DatastreamToPostgresDML dml = + (DatastreamToPostgresDML) DatastreamToPostgresDML.of(null).withColumnCasing("UPPERCASE"); + String json = + "{\"_metadata_primary_keys\": [\"ROWID\"], \"_metadata_deleted\": true," + + " \"_metadata_row_id\": \"AAAEARAAEAAAAC9AAS\"}"; + JsonNode rowObj = getRowObj(json); + List primaryKeys = Arrays.asList("ROWID"); + Map tableSchema = new HashMap<>(); + tableSchema.put("ROWID", "VARCHAR"); + + String filterSql = dml.getPrimaryKeyToValueFilterSql(rowObj, primaryKeys, tableSchema); + assertEquals("\"ROWID\"='AAAEARAAEAAAAC9AAS'", filterSql); + } + /** * Tests that {@link DatastreamToDML#getPrimaryKeyToValueFilterSql} generates a rowid filter even * when destination primaryKeys list is empty (table without PK). @@ -1193,39 +1231,6 @@ public void testGetValueSql_supportsCasedRowIdFallback() { assertEquals("'AAAEARAAEAAAAC9AAS'", valueSql); } - /** - * Tests that {@link DatastreamToDML#getValueSql} resolves _metadata_row_id when column is - * _metadata_row_id and record contains rowid. - */ - @Test - public void testGetValueSql_supportsMetadataRowIdFallbackFromRowId() { - DatastreamToPostgresDML dml = DatastreamToPostgresDML.of(null); - String json = "{\"rowid\": \"AAAEARAAEAAAAC9AAS\"}"; - JsonNode rowObj = getRowObj(json); - Map tableSchema = new HashMap<>(); - tableSchema.put("_metadata_row_id", "VARCHAR"); - - String valueSql = dml.getValueSql(rowObj, "_metadata_row_id", tableSchema); - assertEquals("'AAAEARAAEAAAAC9AAS'", valueSql); - } - - /** - * Tests that {@link DatastreamToDML#getValueSql} resolves _metadata_row_id when column is - * _metadata_row_id and record contains cased ROWID. - */ - @Test - public void testGetValueSql_supportsMetadataRowIdFallbackFromCasedRowId() { - DatastreamToPostgresDML dml = - (DatastreamToPostgresDML) DatastreamToPostgresDML.of(null).withColumnCasing("UPPERCASE"); - String json = "{\"ROWID\": \"AAAEARAAEAAAAC9AAS\"}"; - JsonNode rowObj = getRowObj(json); - Map tableSchema = new HashMap<>(); - tableSchema.put("_metadata_row_id", "VARCHAR"); - - String valueSql = dml.getValueSql(rowObj, "_metadata_row_id", tableSchema); - assertEquals("'AAAEARAAEAAAAC9AAS'", valueSql); - } - /** * Tests that {@link DatastreamToDML#getValueSql} warns and returns empty string when requested * column is missing and not a rowid. From c2bc81b284c015407c979b21cd3fe376d046c238 Mon Sep 17 00:00:00 2001 From: Yonatan T Date: Thu, 24 Sep 2026 04:45:50 +0000 Subject: [PATCH 09/10] refactor(datastream-to-sql): simplify single rowid pk resolution in getPrimaryKeyToValueFilterSql --- .../cloud/teleport/v2/utils/DatastreamToDML.java | 14 ++++++-------- 1 file changed, 6 insertions(+), 8 deletions(-) diff --git a/v2/datastream-to-sql/src/main/java/com/google/cloud/teleport/v2/utils/DatastreamToDML.java b/v2/datastream-to-sql/src/main/java/com/google/cloud/teleport/v2/utils/DatastreamToDML.java index 05a6a2cb5e..7e5a6d56d5 100644 --- a/v2/datastream-to-sql/src/main/java/com/google/cloud/teleport/v2/utils/DatastreamToDML.java +++ b/v2/datastream-to-sql/src/main/java/com/google/cloud/teleport/v2/utils/DatastreamToDML.java @@ -551,14 +551,12 @@ public String getPrimaryKeyToValueFilterSql( } } - if (pkToValueSql.isEmpty() && rowObj.has("_metadata_row_id")) { - for (String pk : primaryKeys) { - if (pk.equalsIgnoreCase(this.rowIdColumnName)) { - String columnValue = getValueSql(rowObj, pk, tableSchema); - pkToValueSql = quote(pk) + "=" + columnValue; - break; - } - } + if (pkToValueSql.isEmpty() + && primaryKeys.size() == 1 + && primaryKeys.get(0).equalsIgnoreCase(this.rowIdColumnName) + && rowObj.has("_metadata_row_id")) { + String pk = primaryKeys.get(0); + pkToValueSql = quote(pk) + "=" + getValueSql(rowObj, pk, tableSchema); } } From 1d20004dea35f9555284d80c994c65810f1211ce Mon Sep 17 00:00:00 2001 From: Yonatan T Date: Tue, 29 Sep 2026 06:22:59 +0000 Subject: [PATCH 10/10] fix(datastream-to-sql): simplify rowid fallback for unkeyed delete events and ignore oracle it in ci --- .../v2/templates/DataStreamToBigQueryIT.java | 1 + .../teleport/v2/utils/DatastreamToDML.java | 121 ++++++------------ 2 files changed, 42 insertions(+), 80 deletions(-) diff --git a/v2/datastream-to-bigquery/src/test/java/com/google/cloud/teleport/v2/templates/DataStreamToBigQueryIT.java b/v2/datastream-to-bigquery/src/test/java/com/google/cloud/teleport/v2/templates/DataStreamToBigQueryIT.java index b0bf3f9306..aacce8e027 100644 --- a/v2/datastream-to-bigquery/src/test/java/com/google/cloud/teleport/v2/templates/DataStreamToBigQueryIT.java +++ b/v2/datastream-to-bigquery/src/test/java/com/google/cloud/teleport/v2/templates/DataStreamToBigQueryIT.java @@ -189,6 +189,7 @@ public void testDataStreamMySqlToBigQueryJson() throws IOException { } @Test + @Ignore("Consolidate feature matrix for expensive tests") public void testDataStreamOracleToBigQueryJson() throws IOException { // Run a simple IT simpleJdbcToBigQueryTest( diff --git a/v2/datastream-to-sql/src/main/java/com/google/cloud/teleport/v2/utils/DatastreamToDML.java b/v2/datastream-to-sql/src/main/java/com/google/cloud/teleport/v2/utils/DatastreamToDML.java index 7e5a6d56d5..e072c6d711 100644 --- a/v2/datastream-to-sql/src/main/java/com/google/cloud/teleport/v2/utils/DatastreamToDML.java +++ b/v2/datastream-to-sql/src/main/java/com/google/cloud/teleport/v2/utils/DatastreamToDML.java @@ -36,7 +36,6 @@ import java.util.List; import java.util.Map; import java.util.Objects; -import java.util.stream.Stream; import javax.sql.DataSource; import org.apache.beam.sdk.transforms.DoFn; import org.apache.beam.sdk.values.KV; @@ -376,14 +375,30 @@ public Map getSqlTemplateValues( return sqlTemplateValues; } - public String getValueSql(JsonNode rowObj, String columnName, Map tableSchema) { - if (columnName.equalsIgnoreCase(this.rowIdColumnName) - && !rowObj.has(columnName) - && rowObj.has("_metadata_row_id")) { - columnName = "_metadata_row_id"; + private JsonNode getFieldFromRowObj(JsonNode rowObj, String columnName) { + JsonNode columnObj = rowObj.get(columnName); + if (columnObj != null) { + return columnObj; + } + if (columnName.equalsIgnoreCase(this.rowIdColumnName)) { + if (rowObj.has("_metadata_row_id")) { + return rowObj.get("_metadata_row_id"); + } + if (rowObj.has(this.rowIdColumnName)) { + return rowObj.get(this.rowIdColumnName); + } } + for (Iterator it = rowObj.fieldNames(); it.hasNext(); ) { + String fieldName = it.next(); + if (applyCasingLogic(fieldName, this.columnCasing).equals(columnName)) { + return rowObj.get(fieldName); + } + } + return null; + } - JsonNode columnObj = rowObj.get(columnName); + public String getValueSql(JsonNode rowObj, String columnName, Map tableSchema) { + JsonNode columnObj = getFieldFromRowObj(rowObj, columnName); if (columnObj == null) { LOG.warn("Missing Required Value: {} in {}", columnName, rowObj.toString()); return ""; @@ -498,91 +513,37 @@ public String getColumnsUpdateSql(JsonNode rowObj, Map tableSche return onUpdateSql; } - private List getSourcePrimaryKeys(JsonNode rowObj) { - if (rowObj == null) { - return Collections.emptyList(); - } - JsonNode pksNode = rowObj.get("_metadata_primary_keys"); - if (pksNode != null && pksNode.isArray()) { - List primaryKeys = new ArrayList<>(); - for (JsonNode node : pksNode) { - primaryKeys.add(node.asText()); - } - return primaryKeys; - } - return Collections.emptyList(); - } - public String getPrimaryKeyToValueFilterSql( JsonNode rowObj, List primaryKeys, Map tableSchema) { - List sourcePrimaryKeys = getSourcePrimaryKeys(rowObj); - String pkToValueSql = ""; - - for (String sourcePkName : sourcePrimaryKeys) { - String destinationPkName = applyCasingLogic(sourcePkName, this.columnCasing); - - if (primaryKeys.contains(destinationPkName)) { - String columnValue = getValueSql(rowObj, sourcePkName, tableSchema); - String quotedDestinationPkName = quote(destinationPkName); - - if (pkToValueSql.isEmpty()) { - pkToValueSql = quotedDestinationPkName + "=" + columnValue; - } else { - pkToValueSql = pkToValueSql + " AND " + quotedDestinationPkName + "=" + columnValue; - } + List targetPrimaryKeys = primaryKeys; + if (targetPrimaryKeys.isEmpty()) { + boolean isDelete = + rowObj.has("_metadata_deleted") && rowObj.get("_metadata_deleted").asBoolean(); + if (!isDelete) { + return ""; } - } - - if (pkToValueSql.isEmpty()) { - for (Iterator it = rowObj.fieldNames(); it.hasNext(); ) { - String sourceFieldName = it.next(); - String destinationPkName = applyCasingLogic(sourceFieldName, this.columnCasing); - - if (primaryKeys.contains(destinationPkName)) { - String columnValue = getValueSql(rowObj, sourceFieldName, tableSchema); - String quotedDestinationPkName = quote(destinationPkName); - - if (pkToValueSql.isEmpty()) { - pkToValueSql = quotedDestinationPkName + "=" + columnValue; - } else { - pkToValueSql = pkToValueSql + " AND " + quotedDestinationPkName + "=" + columnValue; - } - } - } - - if (pkToValueSql.isEmpty() - && primaryKeys.size() == 1 - && primaryKeys.get(0).equalsIgnoreCase(this.rowIdColumnName) - && rowObj.has("_metadata_row_id")) { - String pk = primaryKeys.get(0); - pkToValueSql = quote(pk) + "=" + getValueSql(rowObj, pk, tableSchema); - } - } - - boolean isDelete = - rowObj.has("_metadata_deleted") && rowObj.get("_metadata_deleted").asBoolean(); - if (isDelete && pkToValueSql.isEmpty() && primaryKeys.isEmpty()) { String casedRowId = applyCasingLogic(this.rowIdColumnName, this.columnCasing); - if (!tableSchema.containsKey(casedRowId)) { + boolean schemaHasRowId = + tableSchema.containsKey(casedRowId) + || tableSchema.keySet().stream().anyMatch(k -> k.equalsIgnoreCase(casedRowId)); + if (!schemaHasRowId) { throw new DeletedWithoutPrimaryKey( String.format( "Cannot replicate DELETE for table without primary keys: target schema lacks '%s' column.", casedRowId)); } - String sourceRowIdField = - Stream.of(casedRowId, this.rowIdColumnName, "_metadata_row_id") - .filter(rowObj::has) - .findFirst() - .orElse(null); - - if (sourceRowIdField != null) { - String columnValue = getValueSql(rowObj, sourceRowIdField, tableSchema); - pkToValueSql = quote(casedRowId) + "=" + columnValue; - } + targetPrimaryKeys = Collections.singletonList(casedRowId); } - return pkToValueSql; + List filters = new ArrayList<>(); + for (String destPk : targetPrimaryKeys) { + String columnValue = getValueSql(rowObj, destPk, tableSchema); + if (!columnValue.isEmpty()) { + filters.add(quote(destPk) + "=" + columnValue); + } + } + return String.join(" AND ", filters); } private static Connection getConnection(