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 c4b08cefd5..ea23bddf04 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; @@ -81,7 +82,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); } @@ -269,6 +270,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(); } } @@ -320,10 +324,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(); @@ -361,18 +375,36 @@ public Map getSqlTemplateValues( return sqlTemplateValues; } - public String getValueSql(JsonNode rowObj, String columnName, Map tableSchema) { - String columnValue; + 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; + } + + 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 ""; } - if (columnObj.isTextual()) { - columnValue = "\'" + cleanSql(columnObj.textValue()) + "\'"; - } else { - columnValue = columnObj.toString(); - } + String columnValue = + columnObj.isTextual() ? "'" + cleanSql(columnObj.textValue()) + "'" : columnObj.toString(); return cleanDataTypeValueSql( columnValue, applyCasingLogic(columnName, this.columnCasing), tableSchema); } @@ -485,25 +517,34 @@ public String getColumnsUpdateSql(JsonNode rowObj, Map tableSche public String getPrimaryKeyToValueFilterSql( JsonNode rowObj, List primaryKeys, Map tableSchema) { - DatastreamRow row = DatastreamRow.of(rowObj); - List sourcePrimaryKeys = row.getPrimaryKeys(); - 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); + List targetPrimaryKeys = primaryKeys; + if (targetPrimaryKeys.isEmpty()) { + boolean isDelete = + rowObj.has("_metadata_deleted") && rowObj.get("_metadata_deleted").asBoolean(); + if (!isDelete) { + return ""; + } + String casedRowId = applyCasingLogic(this.rowIdColumnName, this.columnCasing); + 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)); + } + targetPrimaryKeys = Collections.singletonList(casedRowId); + } - if (pkToValueSql.isEmpty()) { - pkToValueSql = quotedDestinationPkName + "=" + columnValue; - } else { - pkToValueSql = pkToValueSql + " AND " + quotedDestinationPkName + "=" + columnValue; - } + List filters = new ArrayList<>(); + for (String destPk : targetPrimaryKeys) { + String columnValue = getValueSql(rowObj, destPk, tableSchema); + if (!columnValue.isEmpty()) { + filters.add(quote(destPk) + "=" + columnValue); } } - return pkToValueSql; + return String.join(" AND ", filters); } private static Connection getConnection( 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 d716464769..4d5d85274e 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,231 @@ 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} 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). + */ + @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 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. + */ + @Test + public void testGetPrimaryKeyToValueFilterSql_usesCasedRowIdFallbackWhenNoPK() { + DatastreamToPostgresDML dml = + (DatastreamToPostgresDML) DatastreamToPostgresDML.of(null).withColumnCasing("UPPERCASE"); + String json = "{\"_metadata_deleted\": true, \"_metadata_row_id\": \"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#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 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.