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
146 changes: 140 additions & 6 deletions cli/api/dbadapters/execution_sql.ts
Original file line number Diff line number Diff line change
Expand Up @@ -153,7 +153,7 @@ from (${query}) as insertions`;
case dataform.OnSchemaChange.EXTEND:
case dataform.OnSchemaChange.SYNCHRONIZE:
this.buildIncrementalSchemaChangeTasks(tasks, table);
// Fall through to run the static DML after the procedure alters the schema
break;
case dataform.OnSchemaChange.IGNORE:
const columns = tableMetadata?.fields.map((f) => f.name) || [];
tasks.add(Task.statement(this.getIncrementalDmlStatement(table, columns)));
Expand Down Expand Up @@ -233,12 +233,14 @@ from (${query}) as insertions`;
...table.target,
name: `${table.target.name}_df_temp_${uniqueId}_empty`,
};
const tempTableName = `${table.target.name}_df_temp_${uniqueId}_temp`;

const procedureName = this.createProcedureName(table.target, uniqueId);
const procedureBody = this.incrementalSchemaChangeBody(
table,
this.resolveTarget(table.target),
emptyTempTableTarget,
tempTableName,
);

const createProcedureSql = `CREATE OR REPLACE PROCEDURE ${procedureName}()
Expand Down Expand Up @@ -274,10 +276,23 @@ END;
DROP PROCEDURE IF EXISTS ${procedureName};`;
}

private declareSchemaChangeVariablesSql(onSchemaChange: dataform.OnSchemaChange): string {
private declareSchemaChangeVariablesSql(table: dataform.ITable): string {
const onSchemaChange = table.onSchemaChange || dataform.OnSchemaChange.IGNORE;
const isMerge =
table.incrementalStrategy !== dataform.IncrementalStrategy.INSERT_OVERWRITE &&
table.uniqueKey &&
table.uniqueKey.length > 0;

let sql = `
-- Declare variables for schema comparison and strategy execution.
DECLARE dataform_columns ARRAY<STRING>;
DECLARE dataform_columns_list STRING;`;

if (isMerge) {
sql += `\nDECLARE dataform_columns_merge STRING;`;
}

sql += `
DECLARE temp_table_columns ARRAY<STRUCT<column_name STRING, data_type STRING>>;
DECLARE columns_added ARRAY<STRUCT<column_name STRING, data_type STRING>>;
DECLARE columns_removed ARRAY<STRING>;`;
Expand Down Expand Up @@ -401,27 +416,146 @@ ${this.alterTableAddColumnsSql(qualifiedTargetTableName)}
END IF;`;
}

private cleanupSql(emptyTempTableName: string): string {
private escapeSqlString(str: string): string {
return str.replace(/\\/g, "\\\\").replace(/"/g, '\\"');
}

private executeDynamicDmlSql(
table: dataform.ITable,
qualifiedTargetTableName: string,
tempTableName: string,
query: string,
): string {
const isMerge =
table.incrementalStrategy !== dataform.IncrementalStrategy.INSERT_OVERWRITE &&
table.uniqueKey &&
table.uniqueKey.length > 0;

let sql = `
-- Prepare dynamic column lists and staging table for DML.
SET dataform_columns_list = (
SELECT STRING_AGG(FORMAT("\`%s\`", column_info.column_name), ", ")
FROM UNNEST(temp_table_columns) AS column_info
);`;

if (isMerge) {
sql += `
SET dataform_columns_merge = (
SELECT STRING_AGG(FORMAT("\`%s\` = DATAFORM_SOURCE.\`%s\`", column_info.column_name, column_info.column_name), ", ")
FROM UNNEST(temp_table_columns) AS column_info
);`;
}

sql += `

CREATE OR REPLACE TEMP TABLE \`${tempTableName}\` AS (
${query}
);
`;

switch (table.incrementalStrategy) {
case dataform.IncrementalStrategy.INSERT_OVERWRITE: {
const partitionBy = table.bigquery && table.bigquery.partitionBy;
const updatePartitionFilter = table.bigquery && table.bigquery.updatePartitionFilter;
const incrementalPredicates = table.bigquery && table.bigquery.incrementalPredicates;
const incrementalPredicatesString =
this.buildIncrementalPredicatesString(incrementalPredicates);
const notMatchedBySourceCondition = [
`${partitionBy} IN UNNEST(partitions_for_replacement)`,
updatePartitionFilter ? `and DATAFORM_DEST.${updatePartitionFilter}` : "",
incrementalPredicatesString,
]
.filter(Boolean)
.join(" ");

sql += `
BEGIN
DECLARE partitions_for_replacement DEFAULT (
ARRAY(
SELECT DISTINCT ${partitionBy}
FROM \`${tempTableName}\`
WHERE ${partitionBy} IS NOT NULL
)
);

EXECUTE IMMEDIATE (
"MERGE ${this.escapeSqlString(qualifiedTargetTableName)} DATAFORM_DEST " ||
"USING \`${tempTableName}\` DATAFORM_SOURCE " ||
"ON FALSE " ||
"WHEN NOT MATCHED BY SOURCE AND ${this.escapeSqlString(notMatchedBySourceCondition)} THEN " ||
"DELETE " ||
"WHEN NOT MATCHED BY TARGET THEN " ||
"INSERT (" || dataform_columns_list || ") VALUES (" || dataform_columns_list || ")"
);
END;`;
break;
}
case dataform.IncrementalStrategy.MERGE:
default: {
if (isMerge) {
const updatePartitionFilter = table.bigquery && table.bigquery.updatePartitionFilter;
const incrementalPredicates = table.bigquery && table.bigquery.incrementalPredicates;
const incrementalPredicatesString =
this.buildIncrementalPredicatesString(incrementalPredicates);
const onCondition = [
table.uniqueKey
.map(
(uniqueKeyCol) => `DATAFORM_DEST.${uniqueKeyCol} = DATAFORM_SOURCE.${uniqueKeyCol}`,
)
.join(" and "),
updatePartitionFilter ? `and DATAFORM_DEST.${updatePartitionFilter}` : "",
incrementalPredicatesString,
]
.filter(Boolean)
.join(" ");

sql += `
EXECUTE IMMEDIATE (
"MERGE ${this.escapeSqlString(qualifiedTargetTableName)} DATAFORM_DEST " ||
"USING \`${tempTableName}\` DATAFORM_SOURCE " ||
"ON ${this.escapeSqlString(onCondition)} " ||
"WHEN MATCHED THEN " ||
"UPDATE SET " || dataform_columns_merge || " " ||
"WHEN NOT MATCHED THEN " ||
"INSERT (" || dataform_columns_list || ") VALUES (" || dataform_columns_list || ")"
);`;
} else {
sql += `
EXECUTE IMMEDIATE (
"INSERT INTO ${this.escapeSqlString(qualifiedTargetTableName)} (" || dataform_columns_list || ") " ||
"SELECT " || dataform_columns_list || " FROM \`${tempTableName}\`"
);`;
}
break;
}
}

return sql;
}

private cleanupSql(emptyTempTableName: string, tempTableName: string): string {
return `
-- Cleanup temporary tables.
DROP TABLE IF EXISTS ${emptyTempTableName};
DROP TABLE IF EXISTS \`${tempTableName}\`;
`;
}

private incrementalSchemaChangeBody(
table: dataform.ITable,
qualifiedTargetTableName: string,
emptyTempTableTarget: dataform.ITarget,
tempTableName: string,
): string {
const emptyTempTableName = this.resolveTarget(emptyTempTableTarget);
const query = this.getIncrementalQuery(table);
const onSchemaChange = table.onSchemaChange || dataform.OnSchemaChange.IGNORE;
const statements: string[] = [
this.declareSchemaChangeVariablesSql(onSchemaChange),
this.declareSchemaChangeVariablesSql(table),
this.createEmptyTempTableSql(emptyTempTableName, query),
this.compareSchemasSql(table.target, emptyTempTableTarget),
this.applySchemaChangeStrategySql(table, qualifiedTargetTableName),
this.cleanupSql(emptyTempTableName),
this.executeDynamicDmlSql(table, qualifiedTargetTableName, tempTableName, query),
this.cleanupSql(emptyTempTableName, tempTableName),
];

return statements.join("\n\n");
Expand Down
58 changes: 58 additions & 0 deletions cli/api/execution_sql_test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -178,6 +178,64 @@ suite("ExecutionSql with 'onSchemaChange'", () => {
}
}
});

test("executes dynamic DML inside procedure and does not emit static DML outside procedure for onSchemaChange", () => {
for (const strategy of [
dataform.OnSchemaChange.FAIL,
dataform.OnSchemaChange.EXTEND,
dataform.OnSchemaChange.SYNCHRONIZE,
]) {
const table = {
...baseTable,
onSchemaChange: strategy,
uniqueKey: ["id"],
};
const tasks = executionSql.publishTasks(table, { fullRefresh: false }, tableMetadata);
const builtTasks = tasks.build();
expect(builtTasks.length).to.equal(1);

const fullSql = builtTasks[0].statement;
const [procedurePart, ...rest] = fullSql.split("\nEND;\n");
const afterProcedurePart = rest.join("\nEND;\n");
expect(procedurePart).to.include("SET dataform_columns_list =");
expect(procedurePart).to.include("SET dataform_columns_merge =");
expect(procedurePart).to.include("EXECUTE IMMEDIATE");
expect(afterProcedurePart).to.not.include("insert into");
expect(afterProcedurePart).to.not.include("merge ");
expect(afterProcedurePart).to.not.include("field1");
}
});

test("error handler outside the procedure only drops fully qualified tables", () => {
// The staging table is scoped to the stored procedure, so it is not resolvable
// from the calling script. Referencing it there raises "must be qualified with a
// dataset", which masks the original error and leaks the procedure.
for (const strategy of [
dataform.OnSchemaChange.FAIL,
dataform.OnSchemaChange.EXTEND,
dataform.OnSchemaChange.SYNCHRONIZE,
]) {
const table = {
...baseTable,
onSchemaChange: strategy,
uniqueKey: ["id"],
};
const fullSql = executionSql
.publishTasks(table, { fullRefresh: false }, tableMetadata)
.build()[0].statement;

const errorHandler = fullSql.split("EXCEPTION WHEN ERROR THEN")[1];
const droppedTables = [...errorHandler.matchAll(/DROP TABLE IF EXISTS `([^`]+)`/g)].map(
(match) => match[1],
);
for (const droppedTable of droppedTables) {
expect(
droppedTable,
`error handler for ${dataform.OnSchemaChange[strategy]} drops an unqualified table`,
).to.match(/^[^.]+\.[^.]+\.[^.]+$/);
}
}
});
});

suite("ExecutionSql for property graphs", () => {
Expand Down
56 changes: 32 additions & 24 deletions cli/api/goldens/insert_overwrite_extend.sql
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@ BEGIN

-- Declare variables for schema comparison and strategy execution.
DECLARE dataform_columns ARRAY<STRING>;
DECLARE dataform_columns_list STRING;
DECLARE temp_table_columns ARRAY<STRUCT<column_name STRING, data_type STRING>>;
DECLARE columns_added ARRAY<STRUCT<column_name STRING, data_type STRING>>;
DECLARE columns_removed ARRAY<STRING>;
Expand Down Expand Up @@ -60,40 +61,47 @@ END IF;



-- Cleanup temporary tables.
DROP TABLE IF EXISTS `project-id.dataset-id.incremental_on_schema_change_df_temp_test_uuid_empty`;

END;
BEGIN
CALL `project-id.dataset-id.df_osc_test_uuid`();
EXCEPTION WHEN ERROR THEN
DROP TABLE IF EXISTS `project-id.dataset-id.incremental_on_schema_change_df_temp_test_uuid_empty`;
DROP PROCEDURE IF EXISTS `project-id.dataset-id.df_osc_test_uuid`;
RAISE;
END;
DROP PROCEDURE IF EXISTS `project-id.dataset-id.df_osc_test_uuid`;
CREATE OR REPLACE TEMP TABLE `staging_table_temp_test_uuid` AS (
-- Prepare dynamic column lists and staging table for DML.
SET dataform_columns_list = (
SELECT STRING_AGG(FORMAT("`%s`", column_info.column_name), ", ")
FROM UNNEST(temp_table_columns) AS column_info
);

CREATE OR REPLACE TEMP TABLE `incremental_on_schema_change_df_temp_test_uuid_temp` AS (
select 1 as id, 'a' as field1, 'new' as field2
);

BEGIN
DECLARE partitions_for_replacement DEFAULT (
ARRAY(
SELECT DISTINCT DATE(ts)
FROM `staging_table_temp_test_uuid`
FROM `incremental_on_schema_change_df_temp_test_uuid_temp`
WHERE DATE(ts) IS NOT NULL
)
);

MERGE `project-id.dataset-id.incremental_on_schema_change` DATAFORM_DEST
USING `staging_table_temp_test_uuid` DATAFORM_SOURCE
ON FALSE
WHEN NOT MATCHED BY SOURCE AND DATE(ts) IN UNNEST(partitions_for_replacement)

THEN
DELETE
WHEN NOT MATCHED BY TARGET THEN
INSERT (`id`,`field1`) VALUES (`id`,`field1`);
EXECUTE IMMEDIATE (
"MERGE `project-id.dataset-id.incremental_on_schema_change` DATAFORM_DEST " ||
"USING `incremental_on_schema_change_df_temp_test_uuid_temp` DATAFORM_SOURCE " ||
"ON FALSE " ||
"WHEN NOT MATCHED BY SOURCE AND DATE(ts) IN UNNEST(partitions_for_replacement) THEN " ||
"DELETE " ||
"WHEN NOT MATCHED BY TARGET THEN " ||
"INSERT (" || dataform_columns_list || ") VALUES (" || dataform_columns_list || ")"
);
END;

DROP TABLE IF EXISTS `staging_table_temp_test_uuid`;

-- Cleanup temporary tables.
DROP TABLE IF EXISTS `project-id.dataset-id.incremental_on_schema_change_df_temp_test_uuid_empty`;
DROP TABLE IF EXISTS `incremental_on_schema_change_df_temp_test_uuid_temp`;

END;
BEGIN
CALL `project-id.dataset-id.df_osc_test_uuid`();
EXCEPTION WHEN ERROR THEN
DROP TABLE IF EXISTS `project-id.dataset-id.incremental_on_schema_change_df_temp_test_uuid_empty`;
DROP PROCEDURE IF EXISTS `project-id.dataset-id.df_osc_test_uuid`;
RAISE;
END;
DROP PROCEDURE IF EXISTS `project-id.dataset-id.df_osc_test_uuid`;
22 changes: 18 additions & 4 deletions cli/api/goldens/on_schema_change_extend.sql
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@ BEGIN

-- Declare variables for schema comparison and strategy execution.
DECLARE dataform_columns ARRAY<STRING>;
DECLARE dataform_columns_list STRING;
DECLARE temp_table_columns ARRAY<STRUCT<column_name STRING, data_type STRING>>;
DECLARE columns_added ARRAY<STRUCT<column_name STRING, data_type STRING>>;
DECLARE columns_removed ARRAY<STRING>;
Expand Down Expand Up @@ -60,8 +61,25 @@ END IF;



-- Prepare dynamic column lists and staging table for DML.
SET dataform_columns_list = (
SELECT STRING_AGG(FORMAT("`%s`", column_info.column_name), ", ")
FROM UNNEST(temp_table_columns) AS column_info
);

CREATE OR REPLACE TEMP TABLE `incremental_on_schema_change_df_temp_test_uuid_temp` AS (
select 1 as id, 'a' as field1, 'new' as field2
);

EXECUTE IMMEDIATE (
"INSERT INTO `project-id.dataset-id.incremental_on_schema_change` (" || dataform_columns_list || ") " ||
"SELECT " || dataform_columns_list || " FROM `incremental_on_schema_change_df_temp_test_uuid_temp`"
);


-- Cleanup temporary tables.
DROP TABLE IF EXISTS `project-id.dataset-id.incremental_on_schema_change_df_temp_test_uuid_empty`;
DROP TABLE IF EXISTS `incremental_on_schema_change_df_temp_test_uuid_temp`;

END;
BEGIN
Expand All @@ -72,7 +90,3 @@ EXCEPTION WHEN ERROR THEN
RAISE;
END;
DROP PROCEDURE IF EXISTS `project-id.dataset-id.df_osc_test_uuid`;
insert into `project-id.dataset-id.incremental_on_schema_change`
(`id`,`field1`)
select `id`,`field1`
from (select 1 as id, 'a' as field1, 'new' as field2) as insertions
Loading