Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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(
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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);
}
Expand Down Expand Up @@ -269,6 +270,9 @@ public List<String> getPrimaryKeys(

for (String destPk : destinationPrimaryKeys) {
if (!casedSourceFieldNames.contains(destPk)) {
if (destPk.equalsIgnoreCase(this.rowIdColumnName) && rowObj.has("_metadata_row_id")) {
continue;
}
return this.getDefaultPrimaryKeys();
}
}
Expand Down Expand Up @@ -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<String> 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();
Expand Down Expand Up @@ -361,18 +375,36 @@ public Map<String, String> getSqlTemplateValues(
return sqlTemplateValues;
}

public String getValueSql(JsonNode rowObj, String columnName, Map<String, String> 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<String> 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<String, String> 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);
}
Expand Down Expand Up @@ -485,25 +517,34 @@ public String getColumnsUpdateSql(JsonNode rowObj, Map<String, String> tableSche
public String getPrimaryKeyToValueFilterSql(
JsonNode rowObj, List<String> primaryKeys, Map<String, String> tableSchema) {

DatastreamRow row = DatastreamRow.of(rowObj);
List<String> 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<String> 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<String> 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(
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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<String> 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<String> 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<String> primaryKeys = Arrays.asList("rowid");
Map<String, String> 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<String> primaryKeys = Arrays.asList("rowid");
Map<String, String> 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<String> primaryKeys = Arrays.asList("ROWID");
Map<String, String> 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<String> primaryKeys = Arrays.asList("ROWID");
Map<String, String> 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<String> primaryKeys = java.util.Collections.emptyList();
Map<String, String> 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<String> primaryKeys = java.util.Collections.emptyList();
Map<String, String> 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<String> primaryKeys = java.util.Collections.emptyList();
Map<String, String> 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<String, String> 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<String, String> 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<String> primaryKeys = java.util.Collections.emptyList();
Map<String, String> 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<String> primaryKeys = Arrays.asList("ID", "ORG_ID");
Map<String, String> 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.
Expand Down
Loading