diff --git a/cli/api/dbadapters/execution_sql.ts b/cli/api/dbadapters/execution_sql.ts index b0724cbcd..b513fcac4 100644 --- a/cli/api/dbadapters/execution_sql.ts +++ b/cli/api/dbadapters/execution_sql.ts @@ -19,7 +19,8 @@ export class ExecutionSql { constructor( private readonly project: dataform.IProjectConfig, private readonly dataformCoreVersion: string, - private readonly uniqueIdGenerator: () => string = () => Math.random().toString(36).substring(2) + private readonly uniqueIdGenerator: () => string = () => + Math.random().toString(36).substring(2), ) { this.CompilationSql = new CompilationSql(project, dataformCoreVersion); } @@ -77,7 +78,7 @@ from (${query}) as insertions`; public shouldWriteIncrementally( table: dataform.ITable, runConfig: dataform.IRunConfig, - tableMetadata?: dataform.ITableMetadata + tableMetadata?: dataform.ITableMetadata, ) { return ( (!runConfig.fullRefresh || table.protected) && @@ -89,7 +90,7 @@ from (${query}) as insertions`; public preOps( table: dataform.ITable, runConfig: dataform.IRunConfig, - tableMetadata?: dataform.ITableMetadata + tableMetadata?: dataform.ITableMetadata, ): Task[] { let preOps = table.preOps; if ( @@ -99,13 +100,13 @@ from (${query}) as insertions`; ) { preOps = table.incrementalPreOps; } - return (preOps || []).map(pre => Task.statement(pre)); + return (preOps || []).map((pre) => Task.statement(pre)); } public postOps( table: dataform.ITable, runConfig: dataform.IRunConfig, - tableMetadata?: dataform.ITableMetadata + tableMetadata?: dataform.ITableMetadata, ): Task[] { let postOps = table.postOps; if ( @@ -115,7 +116,7 @@ from (${query}) as insertions`; ) { postOps = table.incrementalPostOps; } - return (postOps || []).map(post => Task.statement(post)); + return (postOps || []).map((post) => Task.statement(post)); } public resolveTarget(target: dataform.ITarget) { @@ -129,16 +130,16 @@ from (${query}) as insertions`; public publishTasks( table: dataform.ITable, runConfig: dataform.IRunConfig, - tableMetadata?: dataform.ITableMetadata + tableMetadata?: dataform.ITableMetadata, ): Tasks { const tasks = new Tasks(); - this.preOps(table, runConfig, tableMetadata).forEach(statement => tasks.add(statement)); + this.preOps(table, runConfig, tableMetadata).forEach((statement) => tasks.add(statement)); const baseTableType = this.baseTableType(table.enumType); if (tableMetadata && tableMetadata.type !== baseTableType) { tasks.add( - Task.statement(this.dropIfExists(table.target, this.oppositeTableType(baseTableType))) + Task.statement(this.dropIfExists(table.target, this.oppositeTableType(baseTableType))), ); } @@ -152,9 +153,9 @@ 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 + // Fall through to run the static DML after the procedure alters the schema case dataform.OnSchemaChange.IGNORE: - const columns = tableMetadata?.fields.map(f => f.name) || []; + const columns = tableMetadata?.fields.map((f) => f.name) || []; tasks.add(Task.statement(this.getIncrementalDmlStatement(table, columns))); break; } @@ -163,7 +164,7 @@ from (${query}) as insertions`; tasks.add(Task.statement(this.createOrReplace(table))); } - this.postOps(table, runConfig, tableMetadata).forEach(statement => tasks.add(statement)); + this.postOps(table, runConfig, tableMetadata).forEach((statement) => tasks.add(statement)); return tasks.concatenate(); } @@ -171,7 +172,7 @@ from (${query}) as insertions`; public createTableTasks( table: dataform.ITable, runConfig: dataform.IRunConfig, - tableMetadata?: dataform.ITableMetadata + tableMetadata?: dataform.ITableMetadata, ): dataform.IExecutionTask[] { return table.disabled ? [] : this.publishTasks(table, runConfig, tableMetadata).build(); } @@ -179,13 +180,13 @@ from (${query}) as insertions`; public createOperationTasks(operation: dataform.IOperation): dataform.IExecutionTask[] { return operation.disabled ? [] - : operation.queries.map(statement => - dataform.ExecutionTask.create({ type: "statement", statement }) + : operation.queries.map((statement) => + dataform.ExecutionTask.create({ type: "statement", statement }), ); } public createPropertyGraphTasks( - propertyGraph: dataform.IPropertyGraph + propertyGraph: dataform.IPropertyGraph, ): dataform.IExecutionTask[] { const statement = `CREATE OR REPLACE PROPERTY GRAPH ${this.resolveTarget(propertyGraph.target)} ` + @@ -215,14 +216,14 @@ from (${query}) as insertions`; return `drop ${this.tableTypeAsSql(type)} if exists ${this.resolveTarget(target)}`; } - private buildIncrementalPredicatesString( - incrementalPredicates?: string[] | null - ): string { - const validPredicates = incrementalPredicates ? incrementalPredicates.filter(p => p.trim() !== "") : []; + private buildIncrementalPredicatesString(incrementalPredicates?: string[] | null): string { + const validPredicates = incrementalPredicates + ? incrementalPredicates.filter((p) => p.trim() !== "") + : []; if (validPredicates.length === 0) { return ""; } - return `and ${validPredicates.map(p => `(${p})`).join(" and ")}`; + return `and ${validPredicates.map((p) => `(${p})`).join(" and ")}`; } private buildIncrementalSchemaChangeTasks(tasks: Tasks, table: dataform.ITable) { @@ -230,14 +231,14 @@ from (${query}) as insertions`; const emptyTempTableTarget = { ...table.target, - name: `${table.target.name}_df_temp_${uniqueId}_empty` + name: `${table.target.name}_df_temp_${uniqueId}_empty`, }; const procedureName = this.createProcedureName(table.target, uniqueId); const procedureBody = this.incrementalSchemaChangeBody( table, this.resolveTarget(table.target), - emptyTempTableTarget + emptyTempTableTarget, ); const createProcedureSql = `CREATE OR REPLACE PROCEDURE ${procedureName}() @@ -248,7 +249,7 @@ END;`; const callProcedureSql = this.safeCallAndDropProcedure( procedureName, - this.resolveTarget(emptyTempTableTarget) + this.resolveTarget(emptyTempTableTarget), ); tasks.add(Task.statement(createProcedureSql)); tasks.add(Task.statement(callProcedureSql)); @@ -257,14 +258,11 @@ END;`; private createProcedureName(target: dataform.ITarget, uniqueId: string): string { return this.resolveTarget({ ...target, - name: `df_osc_${uniqueId}` + name: `df_osc_${uniqueId}`, }); } - private safeCallAndDropProcedure( - procedureName: string, - emptyTempTableName: string - ): string { + private safeCallAndDropProcedure(procedureName: string, emptyTempTableName: string): string { return ` BEGIN CALL ${procedureName}(); @@ -276,6 +274,20 @@ END; DROP PROCEDURE IF EXISTS ${procedureName};`; } + private declareSchemaChangeVariablesSql(onSchemaChange: dataform.OnSchemaChange): string { + let sql = ` +-- Declare variables for schema comparison and strategy execution. +DECLARE dataform_columns ARRAY; +DECLARE temp_table_columns ARRAY>; +DECLARE columns_added ARRAY>; +DECLARE columns_removed ARRAY;`; + + if (onSchemaChange === dataform.OnSchemaChange.SYNCHRONIZE) { + sql += `\nDECLARE invalid_removed_columns ARRAY;`; + } + return sql; + } + private createEmptyTempTableSql(emptyTempTableName: string, query: string): string { return ` -- Create empty table to extract schema of new query. @@ -286,15 +298,10 @@ CREATE OR REPLACE TABLE ${emptyTempTableName} AS ( private compareSchemasSql( target: dataform.ITarget, - emptyTempTableTarget: dataform.ITarget + emptyTempTableTarget: dataform.ITarget, ): string { return ` -- Compare schemas -DECLARE dataform_columns ARRAY; -DECLARE temp_table_columns ARRAY>; -DECLARE columns_added ARRAY>; -DECLARE columns_removed ARRAY; - SET dataform_columns = ( SELECT IFNULL(ARRAY_AGG(DISTINCT column_name), []) FROM \`${target.database}.${target.schema}.INFORMATION_SCHEMA.COLUMNS\` @@ -321,7 +328,7 @@ SET columns_removed = ( private applySchemaChangeStrategySql( table: dataform.ITable, - qualifiedTargetTableName: string + qualifiedTargetTableName: string, ): string { const onSchemaChange = table.onSchemaChange || dataform.OnSchemaChange.IGNORE; let sql = ` @@ -354,7 +361,6 @@ ${this.alterTableAddColumnsSql(qualifiedTargetTableName)} case dataform.OnSchemaChange.SYNCHRONIZE: const uniqueKeys = table.uniqueKey || []; sql += ` -DECLARE invalid_removed_columns ARRAY; SET invalid_removed_columns = ( SELECT IFNULL(ARRAY_AGG(col), []) FROM UNNEST(columns_removed) AS col WHERE col IN UNNEST(${JSON.stringify(uniqueKeys)}) ); @@ -405,18 +411,17 @@ DROP TABLE IF EXISTS ${emptyTempTableName}; private incrementalSchemaChangeBody( table: dataform.ITable, qualifiedTargetTableName: string, - emptyTempTableTarget: dataform.ITarget + emptyTempTableTarget: dataform.ITarget, ): 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.createEmptyTempTableSql(emptyTempTableName, query), - this.compareSchemasSql( - table.target, - emptyTempTableTarget - ), + this.compareSchemasSql(table.target, emptyTempTableTarget), this.applySchemaChangeStrategySql(table, qualifiedTargetTableName), - this.cleanupSql(emptyTempTableName) + this.cleanupSql(emptyTempTableName), ]; return statements.join("\n\n"); @@ -437,7 +442,7 @@ DROP TABLE IF EXISTS ${emptyTempTableName}; } return `create or replace ${table.materialized ? "materialized " : ""}${this.tableTypeAsSql( - this.baseTableType(table.enumType) + this.baseTableType(table.enumType), )} ${this.resolveTarget(table.target)} ${ table.bigquery && table.bigquery.partitionBy ? `partition by ${table.bigquery.partitionBy} ` @@ -459,20 +464,21 @@ DROP TABLE IF EXISTS ${emptyTempTableName}; columns: string[], query: string, uniqueKey: string[], - bigquery: dataform.IBigQueryOptions + bigquery: dataform.IBigQueryOptions, ) { const updatePartitionFilter = bigquery && bigquery.updatePartitionFilter; const incrementalPredicates = bigquery && bigquery.incrementalPredicates; - const incrementalPredicatesString = this.buildIncrementalPredicatesString(incrementalPredicates); - const backtickedColumns = columns.map(column => `\`${column}\``); + const incrementalPredicatesString = + this.buildIncrementalPredicatesString(incrementalPredicates); + const backtickedColumns = columns.map((column) => `\`${column}\``); return ` merge ${this.resolveTarget(target)} DATAFORM_DEST using (${query} ) DATAFORM_SOURCE -on ${uniqueKey.map(uniqueKeyCol => `DATAFORM_DEST.${uniqueKeyCol} = DATAFORM_SOURCE.${uniqueKeyCol}`).join(` and `)} ${updatePartitionFilter ? `and DATAFORM_DEST.${updatePartitionFilter}` : ""} +on ${uniqueKey.map((uniqueKeyCol) => `DATAFORM_DEST.${uniqueKeyCol} = DATAFORM_SOURCE.${uniqueKeyCol}`).join(` and `)} ${updatePartitionFilter ? `and DATAFORM_DEST.${updatePartitionFilter}` : ""} ${incrementalPredicatesString ? ` ${incrementalPredicatesString}` : ""} when matched then - update set ${columns.map(column => `\`${column}\` = DATAFORM_SOURCE.${column}`).join(",")} + update set ${columns.map((column) => `\`${column}\` = DATAFORM_SOURCE.${column}`).join(",")} when not matched then insert (${backtickedColumns.join(",")}) values (${backtickedColumns.join(",")})`; } @@ -481,15 +487,16 @@ when not matched then target: dataform.ITarget, columns: string[], query: string, - bigquery: dataform.IBigQueryOptions + bigquery: dataform.IBigQueryOptions, ): string { const partitionBy = bigquery && bigquery.partitionBy; const updatePartitionFilter = bigquery && bigquery.updatePartitionFilter; const incrementalPredicates = bigquery && bigquery.incrementalPredicates; - const incrementalPredicatesString = this.buildIncrementalPredicatesString(incrementalPredicates); + const incrementalPredicatesString = + this.buildIncrementalPredicatesString(incrementalPredicates); const uniqueId = this.uniqueIdGenerator(); const stagingTableUnqualified = `staging_table_temp_${uniqueId}`; - const backtickedColumns = columns.map(column => `\`${column}\``); + const backtickedColumns = columns.map((column) => `\`${column}\``); const resolveTargetTable = this.resolveTarget(target); return `CREATE OR REPLACE TEMP TABLE \`${stagingTableUnqualified}\` AS ( @@ -519,20 +526,12 @@ END; DROP TABLE IF EXISTS \`${stagingTableUnqualified}\`;`; } - private getIncrementalDmlStatement( - table: dataform.ITable, - columns: string[] - ): string { + private getIncrementalDmlStatement(table: dataform.ITable, columns: string[]): string { const incrementalQuery = this.getIncrementalQuery(table); switch (table.incrementalStrategy) { case dataform.IncrementalStrategy.INSERT_OVERWRITE: - return this.insertOverwrite( - table.target, - columns, - incrementalQuery, - table.bigquery - ); + return this.insertOverwrite(table.target, columns, incrementalQuery, table.bigquery); case dataform.IncrementalStrategy.MERGE: default: if (table.uniqueKey && table.uniqueKey.length > 0) { @@ -541,13 +540,13 @@ DROP TABLE IF EXISTS \`${stagingTableUnqualified}\`;`; columns, incrementalQuery, table.uniqueKey, - table.bigquery + table.bigquery, ); } return this.insertInto( table.target, - columns.map(column => `\`${column}\``), - incrementalQuery + columns.map((column) => `\`${column}\``), + incrementalQuery, ); } } @@ -556,7 +555,7 @@ DROP TABLE IF EXISTS \`${stagingTableUnqualified}\`;`; export function collectEvaluationQueries( queryOrAction: QueryOrAction, concatenate: boolean, - queryModifier: (mod: string) => string = (q: string) => q + queryModifier: (mod: string) => string = (q: string) => q, ): IValidationQuery[] { // TODO: The prefix method (via `queryModifier`) is a bit sketchy. For example after // attaching the `explain` prefix, a table or operation could look like this: @@ -575,37 +574,37 @@ export function collectEvaluationQueries( if (queryOrAction.enumType === dataform.TableType.INCREMENTAL) { const incrementalTableQueries = queryOrAction.incrementalPreOps.concat( queryOrAction.incrementalQuery, - queryOrAction.incrementalPostOps + queryOrAction.incrementalPostOps, ); if (concatenate) { validationQueries.push({ query: concatenateQueries(incrementalTableQueries, queryModifier), - incremental: true + incremental: true, }); } else { - incrementalTableQueries.forEach(q => - validationQueries.push({ query: queryModifier(q), incremental: true }) + incrementalTableQueries.forEach((q) => + validationQueries.push({ query: queryModifier(q), incremental: true }), ); } } const tableQueries = queryOrAction.preOps.concat( queryOrAction.query, - queryOrAction.postOps + queryOrAction.postOps, ); if (concatenate) { validationQueries.push({ - query: concatenateQueries(tableQueries, queryModifier) + query: concatenateQueries(tableQueries, queryModifier), }); } else { - tableQueries.forEach(q => validationQueries.push({ query: queryModifier(q) })); + tableQueries.forEach((q) => validationQueries.push({ query: queryModifier(q) })); } } else if (queryOrAction instanceof dataform.Operation) { if (concatenate) { validationQueries.push({ - query: concatenateQueries(queryOrAction.queries, queryModifier) + query: concatenateQueries(queryOrAction.queries, queryModifier), }); } else { - queryOrAction.queries.forEach(q => validationQueries.push({ query: queryModifier(q) })); + queryOrAction.queries.forEach((q) => validationQueries.push({ query: queryModifier(q) })); } } else if (queryOrAction instanceof dataform.Assertion) { validationQueries.push({ query: queryModifier(queryOrAction.query) }); @@ -617,6 +616,6 @@ export function collectEvaluationQueries( } } return validationQueries - .map(validationQuery => ({ query: validationQuery.query.trim(), ...validationQuery })) - .filter(validationQuery => !!validationQuery.query); + .map((validationQuery) => ({ query: validationQuery.query.trim(), ...validationQuery })) + .filter((validationQuery) => !!validationQuery.query); } diff --git a/cli/api/execution_sql_test.ts b/cli/api/execution_sql_test.ts index 07b1a0e71..2aee473e4 100644 --- a/cli/api/execution_sql_test.ts +++ b/cli/api/execution_sql_test.ts @@ -9,10 +9,10 @@ suite("ExecutionSql with 'onSchemaChange'", () => { const executionSql = new ExecutionSql( { defaultDatabase: "project-id", - defaultSchema: "dataset-id" + defaultSchema: "dataset-id", }, "2.0.0", - () => "test_uuid" + () => "test_uuid", ); const baseTable: dataform.ITable = { @@ -21,10 +21,10 @@ suite("ExecutionSql with 'onSchemaChange'", () => { target: { database: "project-id", schema: "dataset-id", - name: "incremental_on_schema_change" + name: "incremental_on_schema_change", }, query: "select 1 as id, 'a' as field1", - incrementalQuery: "select 1 as id, 'a' as field1, 'new' as field2" + incrementalQuery: "select 1 as id, 'a' as field1, 'new' as field2", }; const tableMetadata: dataform.ITableMetadata = { @@ -32,22 +32,25 @@ suite("ExecutionSql with 'onSchemaChange'", () => { fields: [ { name: "id", - primitive: dataform.Field.Primitive.INTEGER + primitive: dataform.Field.Primitive.INTEGER, }, { name: "field1", - primitive: dataform.Field.Primitive.STRING - } - ] + primitive: dataform.Field.Primitive.STRING, + }, + ], }; test("generates procedure for FAIL strategy", () => { const table = { ...baseTable, - onSchemaChange: dataform.OnSchemaChange.FAIL + onSchemaChange: dataform.OnSchemaChange.FAIL, }; const tasks = executionSql.publishTasks(table, { fullRefresh: false }, tableMetadata); - const procedureSql = tasks.build().map(t => t.statement).join("\n;\n"); + const procedureSql = tasks + .build() + .map((t) => t.statement) + .join("\n;\n"); const expectedSql = fs.readFileSync("cli/api/goldens/on_schema_change_fail.sql", "utf8"); expect(procedureSql).to.equal(expectedSql.trim()); }); @@ -55,10 +58,13 @@ suite("ExecutionSql with 'onSchemaChange'", () => { test("generates procedure for EXTEND strategy", () => { const table = { ...baseTable, - onSchemaChange: dataform.OnSchemaChange.EXTEND + onSchemaChange: dataform.OnSchemaChange.EXTEND, }; const tasks = executionSql.publishTasks(table, { fullRefresh: false }, tableMetadata); - const procedureSql = tasks.build().map(t => t.statement).join("\n;\n"); + const procedureSql = tasks + .build() + .map((t) => t.statement) + .join("\n;\n"); const expectedSql = fs.readFileSync("cli/api/goldens/on_schema_change_extend.sql", "utf8"); expect(procedureSql).to.equal(expectedSql.trim()); }); @@ -67,10 +73,13 @@ suite("ExecutionSql with 'onSchemaChange'", () => { const table = { ...baseTable, onSchemaChange: dataform.OnSchemaChange.SYNCHRONIZE, - uniqueKey: ["id"] + uniqueKey: ["id"], }; const tasks = executionSql.publishTasks(table, { fullRefresh: false }, tableMetadata); - const procedureSql = tasks.build().map(t => t.statement).join("\n;\n"); + const procedureSql = tasks + .build() + .map((t) => t.statement) + .join("\n;\n"); const expectedSql = fs.readFileSync("cli/api/goldens/on_schema_change_synchronize.sql", "utf8"); expect(procedureSql).to.equal(expectedSql.trim()); }); @@ -79,10 +88,13 @@ suite("ExecutionSql with 'onSchemaChange'", () => { const table = { ...baseTable, onSchemaChange: dataform.OnSchemaChange.IGNORE, - uniqueKey: ["id"] + uniqueKey: ["id"], }; const tasks = executionSql.publishTasks(table, { fullRefresh: false }, tableMetadata); - const procedureSql = tasks.build().map(t => t.statement).join("\n;\n"); + const procedureSql = tasks + .build() + .map((t) => t.statement) + .join("\n;\n"); const expectedSql = fs.readFileSync("cli/api/goldens/on_schema_change_ignore.sql", "utf8"); expect(procedureSql).to.equal(expectedSql.trim()); }); @@ -93,11 +105,17 @@ suite("ExecutionSql with 'onSchemaChange'", () => { incrementalStrategy: dataform.IncrementalStrategy.INSERT_OVERWRITE, bigquery: { partitionBy: "DATE(ts)", - incrementalPredicates: ["DATAFORM_DEST.ts >= '2024-01-01'", "DATAFORM_SOURCE.ts >= '2024-01-01'"] - } + incrementalPredicates: [ + "DATAFORM_DEST.ts >= '2024-01-01'", + "DATAFORM_SOURCE.ts >= '2024-01-01'", + ], + }, }; const tasks = executionSql.publishTasks(table, { fullRefresh: false }, tableMetadata); - const sql = tasks.build().map(t => t.statement).join("\n;\n"); + const sql = tasks + .build() + .map((t) => t.statement) + .join("\n;\n"); const expectedSql = fs.readFileSync("cli/api/goldens/insert_overwrite_ignore.sql", "utf8"); expect(sql).to.equal(expectedSql.trim()); }); @@ -108,24 +126,68 @@ suite("ExecutionSql with 'onSchemaChange'", () => { incrementalStrategy: dataform.IncrementalStrategy.INSERT_OVERWRITE, onSchemaChange: dataform.OnSchemaChange.EXTEND, bigquery: { - partitionBy: "DATE(ts)" - } + partitionBy: "DATE(ts)", + }, }; const tasks = executionSql.publishTasks(table, { fullRefresh: false }, tableMetadata); - const sql = tasks.build().map(t => t.statement).join("\n;\n"); + const sql = tasks + .build() + .map((t) => t.statement) + .join("\n;\n"); const expectedSql = fs.readFileSync("cli/api/goldens/insert_overwrite_extend.sql", "utf8"); expect(sql).to.equal(expectedSql.trim()); }); + + test("places all DECLARE statements before any executable statement in generated procedure body", () => { + 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 createProcedureSql = tasks.build()[0].statement; + expect(createProcedureSql).to.include("CREATE OR REPLACE PROCEDURE"); + + const procedureBody = createProcedureSql.split("BEGIN\n")[1].split("\nEND;")[0]; + const statements = procedureBody + .split(";") + .map((s) => + s + .split("\n") + .filter((line) => !line.trim().startsWith("--")) + .join("\n") + .trim(), + ) + .filter((s) => s.length > 0); + + let seenNonDeclare = false; + for (const stmt of statements) { + if (stmt.toUpperCase().startsWith("DECLARE ")) { + expect( + seenNonDeclare, + `DECLARE statement appeared after non-DECLARE statement in strategy ${dataform.OnSchemaChange[strategy]}: "${stmt}"`, + ).to.equal(false); + } else { + seenNonDeclare = true; + } + } + } + }); }); suite("ExecutionSql for property graphs", () => { const executionSql = new ExecutionSql( { defaultDatabase: "project-id", - defaultSchema: "dataset-id" + defaultSchema: "dataset-id", }, "2.0.0", - () => "test_uuid" + () => "test_uuid", ); test("emits CREATE OR REPLACE PROPERTY GRAPH for FinGraph", () => { @@ -139,10 +201,10 @@ EDGE TABLES ( )`; const propertyGraph: dataform.IPropertyGraph = { target: { database: "project-id", schema: "dataset-id", name: "FinGraph" }, - graphBody + graphBody, }; const tasks = executionSql.createPropertyGraphTasks(propertyGraph); - const sql = tasks.map(t => t.statement).join("\n;\n"); + const sql = tasks.map((t) => t.statement).join("\n;\n"); const expectedSql = fs.readFileSync("cli/api/goldens/property_graph_fingraph.sql", "utf8"); expect(sql).to.equal(expectedSql.trim()); }); @@ -157,10 +219,10 @@ EDGE TABLES ( )`; const propertyGraph: dataform.IPropertyGraph = { target: { database: "project-id", schema: "dataset-id", name: "HRGraph" }, - graphBody + graphBody, }; const tasks = executionSql.createPropertyGraphTasks(propertyGraph); - const sql = tasks.map(t => t.statement).join("\n;\n"); + const sql = tasks.map((t) => t.statement).join("\n;\n"); const expectedSql = fs.readFileSync("cli/api/goldens/property_graph_hrgraph.sql", "utf8"); expect(sql).to.equal(expectedSql.trim()); }); diff --git a/cli/api/goldens/insert_overwrite_extend.sql b/cli/api/goldens/insert_overwrite_extend.sql index 112d475e1..89e6c7255 100644 --- a/cli/api/goldens/insert_overwrite_extend.sql +++ b/cli/api/goldens/insert_overwrite_extend.sql @@ -2,6 +2,13 @@ CREATE OR REPLACE PROCEDURE `project-id.dataset-id.df_osc_test_uuid`() OPTIONS(strict_mode=false) BEGIN +-- Declare variables for schema comparison and strategy execution. +DECLARE dataform_columns ARRAY; +DECLARE temp_table_columns ARRAY>; +DECLARE columns_added ARRAY>; +DECLARE columns_removed ARRAY; + + -- Create empty table to extract schema of new query. CREATE OR REPLACE TABLE `project-id.dataset-id.incremental_on_schema_change_df_temp_test_uuid_empty` AS ( SELECT * FROM (select 1 as id, 'a' as field1, 'new' as field2) AS insertions LIMIT 0 @@ -9,11 +16,6 @@ CREATE OR REPLACE TABLE `project-id.dataset-id.incremental_on_schema_change_df_t -- Compare schemas -DECLARE dataform_columns ARRAY; -DECLARE temp_table_columns ARRAY>; -DECLARE columns_added ARRAY>; -DECLARE columns_removed ARRAY; - SET dataform_columns = ( SELECT IFNULL(ARRAY_AGG(DISTINCT column_name), []) FROM `project-id.dataset-id.INFORMATION_SCHEMA.COLUMNS` diff --git a/cli/api/goldens/on_schema_change_extend.sql b/cli/api/goldens/on_schema_change_extend.sql index de070956e..0b6b520d0 100644 --- a/cli/api/goldens/on_schema_change_extend.sql +++ b/cli/api/goldens/on_schema_change_extend.sql @@ -2,6 +2,13 @@ CREATE OR REPLACE PROCEDURE `project-id.dataset-id.df_osc_test_uuid`() OPTIONS(strict_mode=false) BEGIN +-- Declare variables for schema comparison and strategy execution. +DECLARE dataform_columns ARRAY; +DECLARE temp_table_columns ARRAY>; +DECLARE columns_added ARRAY>; +DECLARE columns_removed ARRAY; + + -- Create empty table to extract schema of new query. CREATE OR REPLACE TABLE `project-id.dataset-id.incremental_on_schema_change_df_temp_test_uuid_empty` AS ( SELECT * FROM (select 1 as id, 'a' as field1, 'new' as field2) AS insertions LIMIT 0 @@ -9,11 +16,6 @@ CREATE OR REPLACE TABLE `project-id.dataset-id.incremental_on_schema_change_df_t -- Compare schemas -DECLARE dataform_columns ARRAY; -DECLARE temp_table_columns ARRAY>; -DECLARE columns_added ARRAY>; -DECLARE columns_removed ARRAY; - SET dataform_columns = ( SELECT IFNULL(ARRAY_AGG(DISTINCT column_name), []) FROM `project-id.dataset-id.INFORMATION_SCHEMA.COLUMNS` diff --git a/cli/api/goldens/on_schema_change_fail.sql b/cli/api/goldens/on_schema_change_fail.sql index 85ad5ed22..17196b41a 100644 --- a/cli/api/goldens/on_schema_change_fail.sql +++ b/cli/api/goldens/on_schema_change_fail.sql @@ -2,6 +2,13 @@ CREATE OR REPLACE PROCEDURE `project-id.dataset-id.df_osc_test_uuid`() OPTIONS(strict_mode=false) BEGIN +-- Declare variables for schema comparison and strategy execution. +DECLARE dataform_columns ARRAY; +DECLARE temp_table_columns ARRAY>; +DECLARE columns_added ARRAY>; +DECLARE columns_removed ARRAY; + + -- Create empty table to extract schema of new query. CREATE OR REPLACE TABLE `project-id.dataset-id.incremental_on_schema_change_df_temp_test_uuid_empty` AS ( SELECT * FROM (select 1 as id, 'a' as field1, 'new' as field2) AS insertions LIMIT 0 @@ -9,11 +16,6 @@ CREATE OR REPLACE TABLE `project-id.dataset-id.incremental_on_schema_change_df_t -- Compare schemas -DECLARE dataform_columns ARRAY; -DECLARE temp_table_columns ARRAY>; -DECLARE columns_added ARRAY>; -DECLARE columns_removed ARRAY; - SET dataform_columns = ( SELECT IFNULL(ARRAY_AGG(DISTINCT column_name), []) FROM `project-id.dataset-id.INFORMATION_SCHEMA.COLUMNS` diff --git a/cli/api/goldens/on_schema_change_synchronize.sql b/cli/api/goldens/on_schema_change_synchronize.sql index d3dc1d223..026e9f4dc 100644 --- a/cli/api/goldens/on_schema_change_synchronize.sql +++ b/cli/api/goldens/on_schema_change_synchronize.sql @@ -2,6 +2,14 @@ CREATE OR REPLACE PROCEDURE `project-id.dataset-id.df_osc_test_uuid`() OPTIONS(strict_mode=false) BEGIN +-- Declare variables for schema comparison and strategy execution. +DECLARE dataform_columns ARRAY; +DECLARE temp_table_columns ARRAY>; +DECLARE columns_added ARRAY>; +DECLARE columns_removed ARRAY; +DECLARE invalid_removed_columns ARRAY; + + -- Create empty table to extract schema of new query. CREATE OR REPLACE TABLE `project-id.dataset-id.incremental_on_schema_change_df_temp_test_uuid_empty` AS ( SELECT * FROM (select 1 as id, 'a' as field1, 'new' as field2) AS insertions LIMIT 0 @@ -9,11 +17,6 @@ CREATE OR REPLACE TABLE `project-id.dataset-id.incremental_on_schema_change_df_t -- Compare schemas -DECLARE dataform_columns ARRAY; -DECLARE temp_table_columns ARRAY>; -DECLARE columns_added ARRAY>; -DECLARE columns_removed ARRAY; - SET dataform_columns = ( SELECT IFNULL(ARRAY_AGG(DISTINCT column_name), []) FROM `project-id.dataset-id.INFORMATION_SCHEMA.COLUMNS` @@ -39,7 +42,6 @@ SET columns_removed = ( -- Apply schema change strategy (SYNCHRONIZE). -DECLARE invalid_removed_columns ARRAY; SET invalid_removed_columns = ( SELECT IFNULL(ARRAY_AGG(col), []) FROM UNNEST(columns_removed) AS col WHERE col IN UNNEST(["id"]) );